1. 金融数据服务从零搭建的完整思路
1.1 为什么我要自己搭一套金融数据服务
最早接触金融数据这块,是因为我需要一套能稳定拉取行情、财报、宏观指标的接口层。市面上的商业数据终端一年动辄几万块,对于个人开发者或者小团队来说成本太高;而免费的数据源要么限流严重,要么字段残缺,要么接口说关就关。踩了几次坑之后我决定自己搭一套financial-services,把数据采集、清洗、存储、对外接口这几层全部掌握在自己手里。
这套服务的核心目标很明确:用最低的成本,拿到尽可能结构化、可追溯、可扩展的金融数据。它解决的不是“我要做量化交易”这种终极问题,而是“我需要一个稳定的数据底座,让上层应用不用天天担心数据断供”这个前置问题。适合谁来参考?我认为有三类人:一是做个人量化研究、需要批量历史数据的开发者;二是做金融类应用、需要自建数据中台的工程师;三是对数据管道感兴趣、想拿一个真实场景练手的技术爱好者。
需要提前说明的是,本文讨论的所有数据均来自公开、合规的渠道,涉及具体数据源时我会讲清楚获取方式和注意事项,但不会给出任何绕过授权或违反服务条款的做法。金融数据的合规使用是底线,这一点在后面每个环节我都会反复强调。
1.2 整体架构的分层设计
我把整套financial-services拆成了四层,这个分层不是拍脑袋定的,而是根据数据从“生”到“用”的流转路径倒推出来的。
第一层是采集层,负责从各个数据源把原始数据拉回来。这一层的关键词是“容错”和“限速”。金融数据源普遍有频率限制,粗暴并发只会被封,所以采集层必须内置退避重试和请求队列。
第二层是清洗层,负责把原始数据标准化。不同来源的字段名、时间格式、复权方式都不一样,比如有的用trade_date,有的用date,有的返回时间戳,有的返回字符串。清洗层的任务就是把这些统一成一套内部 schema。
第三层是存储层,负责持久化。我选的是关系型数据库加对象存储的组合:结构化数据(行情、财报)进数据库,原始响应体和日志进对象存储,方便回溯。
第四层是服务层,对外暴露统一的 REST 接口,屏蔽底层数据源的差异。上层应用只认这一套接口,数据源换了、加了,上层无感知。
提示:分层最大的好处是“替换成本低”。某个数据源挂了,我只需要改采集层的一个适配器,清洗、存储、服务三层完全不用动。这一点在长期维护中价值极高。
1.3 技术选型背后的取舍逻辑
选型这块我纠结了很久,最后定下来的组合是 Python + PostgreSQL + Redis + FastAPI。逐个说下理由。
Python几乎是金融数据处理的事实标准,pandas、numpy 生态成熟,各种数据源的 SDK 也最全。虽然性能上不如 Go 或 Rust,但金融数据服务的瓶颈通常在 IO 和限流,不在计算,Python 完全够用。
PostgreSQL而不是 MySQL,主要看中它对 JSON 字段、窗口函数、时序查询的支持更好。金融数据里财报这种半结构化数据用 JSONB 存非常方便,而计算移动平均、同比环比这类指标,窗口函数能省掉大量应用层代码。
Redis用来做接口缓存和限流计数器。行情类数据读多写少,缓存命中率很高,能显著降低数据库压力。
FastAPI是因为它自带异步支持和自动生成的接口文档,对于需要频繁调试接口的数据服务来说,省心。
这里我要强调一个取舍:不要一上来就上分布式和消息队列。我见过太多人搭数据服务第一步就上 Kafka、Celery,结果数据量根本没到那个级别,运维复杂度却翻了好几倍。我的建议是单机跑通全流程,等真的遇到性能瓶颈再拆。
2. 数据采集层的核心细节与实操要点
2.1 数据源分类与适配策略
金融数据源大致可以分成三类,每类的采集策略完全不同。
第一类是行情数据,包括股票、指数、期货的日线、分钟线。这类数据的特点是量大、更新频繁、对时效性要求高。采集策略是增量拉取,每天收盘后拉当天数据,历史数据一次性初始化。
第二类是财务数据,包括财报三张表、财务指标。这类数据更新频率低(季度),但字段多、结构复杂。采集策略是全量拉取加版本对比,因为财报会有更正公告,需要能追溯历史版本。
第三类是宏观与参考数据,包括利率、汇率、行业分类、交易日历。这类数据量小但重要性高,交易日历尤其关键,它决定了整个系统的调度节奏。
针对这三类,我在采集层写了统一的适配器接口,每个数据源实现fetch、parse、normalize三个方法。这样新增数据源只需要写一个适配器类,不用动主流程。
2.2 限速与重试机制的具体实现
限速这块踩过的坑最多。早期我天真地以为只要加个sleep就行,结果发现不同接口的限流规则完全不一样:有的按分钟算,有的按天算,有的对单 IP 限,有的对账号限。
我最终的方案是令牌桶 + 指数退避的组合。令牌桶控制请求速率,指数退避处理被限流后的重试。核心逻辑大概是这样:
import time import random class RateLimiter: def __init__(self, rate_per_sec, burst): self.rate = rate_per_sec self.burst = burst self.tokens = burst self.last = time.time() def acquire(self): now = time.time() self.tokens = min(self.burst, self.tokens + (now - self.last) * self.rate) self.last = now if self.tokens < 1: wait = (1 - self.tokens) / self.rate time.sleep(wait) self.tokens = 0 else: self.tokens -= 1 def fetch_with_retry(func, max_retries=5): for attempt in range(max_retries): try: return func() except RateLimitError: backoff = (2 ** attempt) + random.uniform(0, 1) time.sleep(backoff) raise Exception("max retries exceeded")指数退避里加随机抖动(jitter)很关键,否则多个任务同时被限流后会同时重试,形成“惊群”,反而更容易再次触发限流。这个细节很多教程不讲,但实际生产里非常有用。
注意:限速参数不要照抄别人的配置。每个数据源的规则不同,一定要先小批量测试,观察响应头和错误码,摸清真实阈值再定参数。我一般会留 20% 的余量,比如接口允许每分钟 100 次,我就配 80 次。
2.3 数据源切换与降级方案
单一数据源是脆弱的,这是我用血泪换来的教训。有一次主力数据源临时维护,整个服务停摆了两天。从那以后我给每个关键数据都配了至少一个备用源。
降级方案的设计要点是字段对齐。备用源的字段名、单位、复权方式可能和主源不同,所以清洗层必须能识别数据来自哪个源,并做相应转换。我的做法是在数据表里加一个source字段,记录每条数据的来源,这样出问题时能快速定位,也方便做数据质量对比。
切换逻辑我放在采集层的调度器里:主源连续失败 N 次,自动切到备用源,同时发告警。N 我设的是 3,太小容易误切,太大又失去意义。切换后不会自动切回,需要人工确认主源恢复,避免来回抖动。
3. 数据清洗与存储的实操过程
3.1 字段标准化与时间处理
清洗层最繁琐的工作是字段标准化。我定义了一套内部 schema,所有数据源的数据都要映射到这套 schema 上。以日线行情为例,内部字段是symbol、trade_date、open、high、low、close、volume、amount,不管数据源怎么叫,最后都统一成这套。
时间处理是另一个重灾区。有的源返回20240101,有的返回2024-01-01,有的返回1704067200时间戳。我统一转成date类型存库,转换函数里对每种格式做正则匹配。这里有个坑:时区。有些源返回的是 UTC 时间,如果不处理,日线数据会错位一天。我的做法是所有时间统一按交易所所在时区处理,存库时只存日期,不存时间戳,从根上避免时区问题。
复权处理也要在清洗层做。前复权、后复权、不复权三种口径,我选择存不复权的原始价格,复权因子单独存一张表,查询时按需计算。这样做的原因是复权因子会随新的除权除息事件变化,如果直接存复权后的价格,历史数据就得反复重算,而存原始价格加因子,历史数据永远不用动。
3.2 数据库表结构设计
表结构设计我遵循一个原则:宽表存快照,窄表存明细。行情这种每天一行、字段固定的数据,用宽表;财报这种字段多、可能增减的,用窄表加 JSONB。
日线行情表大概长这样:
CREATE TABLE daily_quote ( symbol VARCHAR(16) NOT NULL, trade_date DATE NOT NULL, open NUMERIC(18, 4), high NUMERIC(18, 4), low NUMERIC(18, 4), close NUMERIC(18, 4), volume BIGINT, amount NUMERIC(24, 4), source VARCHAR(32), created_at TIMESTAMP DEFAULT NOW(), PRIMARY KEY (symbol, trade_date) );主键用(symbol, trade_date)而不是自增 ID,这样天然去重,重复插入会直接报错,配合ON CONFLICT DO UPDATE就能实现幂等写入。这个设计让我的采集任务可以放心重跑,不用担心产生重复数据。
财报数据我用的是symbol + report_date + field_name + field_value的窄表结构,虽然查询时要 pivot,但字段增减完全不用改表结构,灵活性高很多。
3.3 幂等写入与增量更新
幂等是数据服务的基本功。我的所有写入操作都设计成可重复执行:行情用ON CONFLICT DO UPDATE,财报用先删后插(按symbol + report_date删,再插新数据)。
增量更新的关键是水位线。我维护一张sync_state表,记录每个数据源每个数据集最后同步到的位置。每次采集从水位线往后拉,拉完更新水位线。这样既不会漏数据,也不会重复拉。
CREATE TABLE sync_state ( dataset VARCHAR(64) PRIMARY KEY, last_sync_date DATE, last_sync_at TIMESTAMP, status VARCHAR(16) );水位线更新必须和业务数据写入在同一个事务里,否则可能出现数据写入了但水位线没更新(重复拉)或者水位线更新了但数据没写入(漏数据)的情况。这个事务边界问题我在早期没注意,导致过一次数据缺口,排查了很久才发现。
4. 服务层接口设计与常见问题排查
4.1 REST 接口的字段与分页设计
服务层对外暴露的接口我尽量保持简洁,核心就几个:查行情、查财报、查交易日历、查同步状态。每个接口的返回结构统一成{code, message, data}三段式,方便前端统一处理。
分页这块我踩过一个坑:早期用offset + limit分页,数据量大了之后深分页性能急剧下降。后来改成游标分页,用trade_date或自增 ID 作为游标,性能稳定。对于金融数据这种天然有序的数据,游标分页几乎是必然选择。
@app.get("/api/quote/{symbol}") async def get_quote(symbol: str, start: str, end: str, cursor: str = None, limit: int = 500): query = build_query(symbol, start, end, cursor, limit) rows = await db.fetch_all(query) next_cursor = rows[-1]["trade_date"] if len(rows) == limit else None return {"code": 0, "message": "ok", "data": rows, "next_cursor": next_cursor}缓存策略上,历史数据(比如一个月前的)基本不变,缓存时间长;近期数据缓存时间短。我用 Redis 做两级缓存,key 里带上参数哈希,TTL 按数据新鲜度动态设置。
4.2 常见问题速查表
下面这张表是我运维这套服务一年多积累下来的高频问题,基本覆盖了 80% 的故障场景。
| 问题现象 | 可能原因 | 排查方向 | 解决方法 |
|---|---|---|---|
| 采集任务卡住不动 | 数据源限流或网络超时 | 看日志最后一条请求 | 检查限速配置,加超时和重试 |
| 数据出现缺口 | 水位线更新与写入不在同一事务 | 对比 sync_state 和实际数据 | 修复事务边界,补拉缺口数据 |
| 接口响应变慢 | 深分页或缓存失效 | 看慢查询日志 | 改游标分页,调整缓存 TTL |
| 财报数据对不上 | 复权口径或字段映射错误 | 对比原始响应和入库数据 | 检查清洗层映射规则 |
| 重复数据 | 幂等逻辑失效 | 查主键冲突日志 | 确认 ON CONFLICT 生效 |
| 时区错位 | 时间未统一处理 | 抽查跨时区数据 | 统一按交易所时区转换 |
4.3 监控与告警的落地经验
监控这块我一开始没重视,直到有次数据断了三天才发现。后来我加了三层监控:采集层监控任务成功率,存储层监控数据新鲜度,服务层监控接口延迟和错误率。
数据新鲜度监控最有用。逻辑很简单:查每个数据集最新一条数据的时间,如果超过预期更新周期还没更新,就告警。比如日线数据,如果今天收盘后两小时还没有今天的数据,就说明采集出问题了。这个监控帮我提前发现了好几次数据源异常。
告警渠道我用的是邮件加即时通讯工具的 webhook,分级发送:一般问题发邮件,严重问题(比如连续失败)直接发即时消息。告警一定要克制,否则会被淹没,我给自己定的规则是“只有需要人工介入的才告警”。
提示:监控指标要能回答“数据是不是新鲜的、全的、对的”这三个问题。新鲜看更新时间,全看行数和覆盖标的数,对看抽样校验。三者缺一不可。
5. 我在实操中踩过的坑与经验总结
5.1 关于数据源稳定性的真实体会
做金融数据服务,最大的不确定性永远来自数据源。我总结了一条经验:永远假设数据源会挂。基于这个假设,所有设计都要考虑降级和恢复。
具体来说,采集任务要能断点续传,数据要能补拉,接口要有缓存兜底。我甚至给关键接口做了“最后已知良好值”的兜底逻辑:如果实时查询失败,返回缓存里最近一次成功的数据,并标记stale: true,让上层自己决定要不要用。这个设计在数据源抖动时极大提升了可用性。
另外,数据源的接口文档经常和实际行为不一致,字段含义、返回格式、错误码都可能变。我的做法是每次采集都记录原始响应体的哈希,一旦哈希变化就告警,人工确认是不是接口变了。这个“响应指纹”机制帮我抓到过好几次数据源的静默变更。
5.2 性能优化的几个关键点
性能优化我遵循“先测量再优化”的原则。用 profiling 工具找出真正的瓶颈,而不是凭感觉优化。
实测下来,这套服务的瓶颈主要在三个地方:数据库写入、接口序列化、缓存穿透。数据库写入用批量插入加事务能提升一个数量级;接口序列化用orjson替代标准库能快好几倍;缓存穿透用布隆过滤器或空值缓存解决。
还有一个容易被忽略的点是连接池。数据库连接和 HTTP 连接都要用池化,否则频繁建连的开销会吃掉大量性能。我把连接池大小设成并发数的 1.5 倍左右,实测比较稳。
5.3 后续可以扩展的方向
这套服务目前满足了我的核心需求,但还有不少可以扩展的地方。比如加一层数据质量校验,对入库数据做完整性、一致性、异常值检查;比如加回测数据接口,直接对外提供复权后的连续价格序列;再比如把采集调度从单机 cron 换成更灵活的调度框架,支持依赖编排。
不过我的建议是按需扩展,不要过度设计。每加一个功能都要问自己:现在真的需要吗?我见过太多项目因为过早引入复杂架构而烂尾。先把核心链路跑稳,等真实需求出现再动手,这是我用很多时间换来的教训。
最后分享一个小技巧:给每个数据源写一个“健康检查”脚本,定期跑一遍,输出成功率、延迟、数据量等指标。这个脚本不用很复杂,几十行就够,但能让你对数据源状态心里有数,出问题时第一时间知道是哪个源的问题。