
1. 金融数据服务从零搭建的完整思路1.1 为什么我要自己搭一套金融数据服务先说清楚这个项目到底在干什么。financial-services这个名字听起来很泛但落到实际工程里它指的是一套面向金融场景的数据服务层——把行情、财报、宏观经济指标、汇率、利率这些散落在各处的数据统一采集、清洗、存储再通过标准接口对外提供查询能力。它解决的核心问题是数据源太多太杂格式不统一直接对接业务代码会让整个系统变成一团乱麻。我做这套东西的起因很直接。之前手上有个投研分析的小工具需要同时拉A股行情、美股指数、国内宏观经济数据还要算一些简单的财务比率。一开始图省事每个数据源写一个脚本直接调结果三个月后代码里全是各种 API 的密钥、不同的时间格式、五花八门的字段命名改一个数据源要翻五个文件。后来痛定思痛决定抽一层出来把所有数据源统一收口这就是financial-services的雏形。这套服务适合谁参考如果你正在做量化回测、投研工具、财务分析系统或者任何需要稳定获取金融数据的应用这套架构都能直接拿去用。哪怕你只是想定期抓一些财经数据做个人看板里面的采集调度、数据清洗、缓存策略这些模块也能拆出来单独用。技术栈上我用的是 Python 为主FastAPI 做接口层PostgreSQL 存结构化数据Redis 做缓存整体不依赖什么冷门组件一台普通云主机就能跑起来。1.2 整体架构怎么分层才不乱金融数据服务和普通业务服务最大的区别在于数据源极不稳定且对时效性和准确性要求极高。行情数据可能每秒都在变财报数据一个季度才更新一次宏观数据按月发布。如果用一个统一的采集频率去处理所有数据要么浪费资源要么错过关键更新。所以我的分层思路是这样的最底层是数据源适配层每个数据源一个适配器负责处理该数据源特有的认证、分页、限流、字段映射往上一层是采集调度层用不同的调度策略驱动不同的适配器行情类高频、财报类低频、宏观类按发布时间触发再往上是清洗与标准化层把所有数据统一成内部标准格式比如时间一律用 UTC 时间戳金额一律用最小货币单位字段命名统一用蛇形命名法最上面是服务接口层对外提供 RESTful 接口和批量查询能力。这样分层的好处是新增一个数据源只需要写一个适配器不用动其他任何代码。我后来加港股数据的时候从写适配器到上线测试半天就搞定了。如果当初把所有逻辑揉在一起加一个源至少得改一周。提示分层的时候一定要把“数据源特有的逻辑”和“通用逻辑”严格分开。我见过太多项目把某个数据源的字段名直接透传到接口层结果换个数据源整个前端都得跟着改。1.3 技术选型背后的取舍逻辑选 PostgreSQL 而不是 MySQL主要考虑的是金融数据里大量涉及时间序列查询和复杂聚合。PostgreSQL 的窗口函数、CTE、以及BRIN索引在处理按时间范围扫描的场景下表现更稳。而且JSONB字段类型让我可以在标准化表结构之外保留原始数据的完整快照方便排查问题。缓存用 Redis 是常规操作但这里有个细节不同数据的缓存策略完全不同。实时行情缓存 3 到 5 秒日线数据缓存到当天收盘财报数据缓存 24 小时宏观数据缓存 7 天。我一开始用统一的 60 秒过期结果财报数据被反复拉取白白浪费了 API 配额。后来改成按数据类型配置不同的 TTLAPI 调用量直接降了七成。接口层选 FastAPI一是异步支持好二是自动生成 OpenAPI 文档前端同事对接起来不用我反复解释字段含义。至于调度我没上 Celery 这种重家伙而是用 APScheduler 加一个简单的任务队列对于中小规模的数据服务来说完全够用运维成本也低得多。2. 核心模块的细节拆解与实操要点2.1 数据源适配器怎么写才通用适配器是整个服务的基石写得好不好直接决定了后续扩展的成本。我的做法是定义一个抽象基类把所有数据源共有的行为抽象出来from abc import ABC, abstractmethod from typing import Any class BaseAdapter(ABC): source_name: str rate_limit: int # 每分钟最大请求数 abstractmethod async def fetch_raw(self, **kwargs) - Any: 从数据源拉取原始数据 pass abstractmethod def normalize(self, raw: Any) - list[dict]: 将原始数据转换为内部标准格式 pass async def fetch(self, **kwargs) - list[dict]: raw await self.fetch_raw(**kwargs) return self.normalize(raw)这个基类里fetch_raw负责处理认证、请求、重试这些脏活normalize负责字段映射和类型转换。每个具体的数据源适配器只需要实现这两个方法。我还在基类里加了限流装饰器用令牌桶算法控制请求频率避免触发数据源的封禁策略。字段映射这块有个坑要特别注意不同数据源对同一概念的定义可能完全不同。比如“成交量”有的源给的是股数有的给的是手数有的甚至是金额。我在normalize里强制要求所有适配器把成交量统一成股数如果源数据是手数就乘以 100。这个转换逻辑必须写在适配器里不能留到上层去猜。注意适配器里绝对不要做业务逻辑判断。我见过有人在适配器里根据数据值决定要不要报警结果换个数据源报警逻辑就失效了。适配器只负责“取数据”和“转格式”其他一概不管。2.2 数据清洗与标准化的关键规则数据清洗是金融数据服务里最耗时间但也最不能省的一步。原始数据里常见的脏东西包括缺失值用各种奇怪符号表示--、N/A、null、空字符串、时间格式五花八门有的带时区有的不带、数字里混着千分位逗号和货币符号、字段名大小写不一致。我的清洗规则分三步走。第一步是类型强制转换所有数值字段先转成Decimal类型避免浮点精度问题。金融计算里用float是自找麻烦0.1 加 0.2 不等于 0.3 这种事在财务对账时能让人崩溃。第二步是缺失值统一处理所有缺失值一律转成None在入库时存为NULL查询时由业务层决定怎么展示。第三步是时间标准化所有时间字段统一转成 UTC 时间戳存储展示时再按用户时区转换。这里有个实操心得清洗规则一定要写成可配置的。我一开始把规则硬编码在代码里后来发现某个数据源突然改了字段格式不得不改代码重新部署。现在我把字段映射和清洗规则放在 YAML 配置文件里改规则只需要改配置重启服务不用动代码。# config/adapters/source_a.yaml fields: trade_date: target: trade_date type: date format: %Y%m%d volume: target: volume type: decimal multiplier: 100 # 手转股 amount: target: amount type: decimal strip_chars: ,¥2.3 缓存策略的精细化配置缓存这块我踩过的坑最多值得单独拿出来说。最开始我用的是最简单的“查缓存没有就查库然后写缓存”模式结果遇到两个问题一是缓存击穿某个热点数据过期瞬间大量请求打到数据库二是缓存和数据库不一致数据更新后缓存还是旧的。解决缓存击穿用的是互斥锁加空值缓存。当缓存未命中时不是所有请求都去查库而是先抢一个分布式锁抢到的去查库并写缓存没抢到的等一小会儿再查缓存。对于确实不存在的数据也缓存一个空值标记避免反复查询。缓存一致性方面我采用的是写时更新加过期兜底。数据更新时主动删除对应缓存同时设置一个较短的过期时间作为兜底。金融数据对一致性要求高但也不是所有数据都需要强一致行情数据差几秒可以接受财报数据差几分钟问题也不大。数据类型缓存 TTL更新策略一致性要求实时行情3 秒写时删除最终一致日线行情至当日收盘写时删除最终一致财务报表24 小时写时删除最终一致宏观指标7 天定时刷新弱一致基础信息12 小时写时删除最终一致2.4 接口设计中的分页与批量查询对外接口设计直接影响到调用方的体验。金融数据查询有两个典型场景一是查单只股票的某段时间数据二是批量查多只股票的某个时点数据。这两种场景对接口的要求完全不同。单标的时序查询用时间范围加游标分页。不要用offset/limit因为金融数据在持续写入用 offset 分页会导致数据重复或遗漏。我的做法是返回一个游标下次查询带上这个游标继续往后取。批量查询用标的列表加字段过滤。调用方传一个标的代码列表和需要的字段列表服务端只返回这些字段避免传输大量无用数据。这里有个细节批量查询一定要限制单次请求的标的数量上限我设的是 200 个超过就分批处理。不设上限的话有人传几千个标的进来数据库直接被打爆。# 批量查询接口示例 app.post(/api/v1/quotes/batch) async def batch_quotes( symbols: list[str], fields: list[str] [close, volume], trade_date: str None ): if len(symbols) 200: raise HTTPException(400, 单次查询标的数不能超过200) # ... 查询逻辑3. 完整实操流程与核心环节实现3.1 环境搭建与依赖安装先把基础环境跑起来。我用的 Python 3.11数据库 PostgreSQL 15缓存 Redis 7。操作系统不限Linux 和 macOS 都行Windows 建议用 WSL2。# 创建虚拟环境 python -m venv venv source venv/bin/activate # Windows 用 venv\Scripts\activate # 安装核心依赖 pip install fastapi uvicorn[standard] sqlalchemy asyncpg redis apscheduler pyyaml httpx pydantic数据库初始化这块我建议用 Alembic 做迁移管理不要手动建表。金融数据的表结构后期调整频率很高手动改表迟早会出乱子。pip install alembic alembic init migrations # 修改 alembic.ini 中的 sqlalchemy.url alembic revision --autogenerate -m init tables alembic upgrade head核心表结构我设计了四张主表instruments存标的的基础信息quotes存行情数据financials存财报数据macro_indicators存宏观指标。每张表都按时间做了分区行情表按月分区财报表按季度分区。分区的好处是查询时能自动裁剪掉不相关的分区速度提升非常明显。提示PostgreSQL 的分区表在主键设计上有个坑分区键必须包含在主键里。我一开始用自增 ID 做主键后来改分区时不得不重建表。建议一开始就用(id, trade_date)这样的复合主键。3.2 数据采集调度的配置与启动调度这块我用 APScheduler 的AsyncIOScheduler和 FastAPI 跑在同一个事件循环里省得再维护一个独立的调度进程。from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.interval import IntervalTrigger scheduler AsyncIOScheduler() # 实时行情交易时段每5秒采集一次 scheduler.add_job( collect_realtime_quotes, IntervalTrigger(seconds5), idrealtime_quotes, max_instances1 ) # 日线行情每天收盘后采集 scheduler.add_job( collect_daily_quotes, CronTrigger(hour15, minute30, day_of_weekmon-fri), iddaily_quotes ) # 财报数据每季度发布期密集采集 scheduler.add_job( collect_financials, CronTrigger(hour*/2, day1-30, month1,4,7,10), idfinancials )这里的关键是max_instances1防止上一次采集还没跑完下一次又启动了。金融数据采集有时候会因为网络问题变慢不加这个限制会出现任务堆积。调度启动后我加了一个简单的监控面板用 Redis 记录每个任务的最后执行时间和执行状态。如果某个任务超过预期时间没执行就发告警。这个监控不复杂但非常有用我有次就是因为数据源接口变更导致采集任务静默失败靠这个监控才发现。3.3 数据入库的批量写入优化数据入库的性能直接决定了整个服务的吞吐量。我一开始用 ORM 逐条插入采集 5000 条行情数据要花将近 30 秒后来改成批量插入同样的数据量 2 秒搞定。from sqlalchemy.dialects.postgresql import insert async def bulk_upsert_quotes(session, quotes: list[dict]): stmt insert(Quote).values(quotes) stmt stmt.on_conflict_do_update( index_elements[symbol, trade_date], set_{ close: stmt.excluded.close, volume: stmt.excluded.volume, amount: stmt.excluded.amount, updated_at: func.now() } ) await session.execute(stmt) await session.commit()用ON CONFLICT DO UPDATE实现 upsert这样重复采集不会产生重复数据还能自动更新最新值。批量大小我设的是 1000 条一批太大容易导致单次事务超时太小又体现不出批量优势。这个值可以根据实际数据库配置调整一般 500 到 2000 之间都合理。还有个细节入库前一定要做数据校验。我遇到过数据源返回的行情数据里混着停牌股票的空值如果不校验直接入库后续查询会报错。校验规则包括价格必须大于零、成交量不能为负、交易日期不能是未来时间。校验不通过的数据记录到日志里人工排查。3.4 接口服务的部署与压测服务用 Uvicorn 启动生产环境建议用 Gunicorn 加 Uvicorn worker 的方式充分利用多核 CPU。gunicorn app.main:app \ --workers 4 \ --worker-class uvicorn.workers.UvicornWorker \ --bind 0.0.0.0:8000 \ --timeout 120 \ --access-logfile -部署完成后一定要做压测。我用 Locust 写了个简单的压测脚本模拟 100 个并发用户查询行情数据。第一次压测结果很不理想QPS 只有 200 左右排查发现是数据库连接池太小。SQLAlchemy 默认连接池是 5 个改成 20 个之后 QPS 直接上到 1500。# 数据库连接池配置 engine create_async_engine( DATABASE_URL, pool_size20, max_overflow10, pool_pre_pingTrue, pool_recycle3600 )pool_pre_ping这个参数建议打开它会在每次从连接池取连接时先 ping 一下避免使用到已经断开的连接。金融数据服务经常长时间运行数据库连接被中间件断开是常有的事不加这个参数会时不时报连接错误。4. 常见问题排查与避坑经验实录4.1 数据源接口变更的应对策略数据源接口变更是这个项目里最让人头疼的问题没有之一。我统计了一下平均每个月都会遇到至少一次数据源字段调整或接口地址变更。应对策略的核心是快速发现、快速定位、快速修复。快速发现靠的是数据校验和监控告警。我在入库前加了一层校验如果某次采集的数据里关键字段缺失率超过 10%就触发告警。这样能在数据源变更的第一时间收到通知而不是等业务方反馈数据不对才发现。快速定位靠的是原始数据快照。每次采集的原始响应我都会存一份到对象存储里保留 7 天。出问题时直接对比原始数据和标准化后的数据一眼就能看出是哪个字段的映射出了问题。快速修复靠的是配置化。前面提到的 YAML 配置在这里发挥了关键作用大部分字段变更只需要改配置不用改代码。只有接口地址或认证方式变更才需要动代码这种情况相对较少。问题类型发现方式修复方式平均修复时间字段名变更校验告警改 YAML 配置5 分钟字段格式变更校验告警改 YAML 配置10 分钟接口地址变更采集失败告警改代码配置30 分钟认证方式变更采集失败告警改适配器代码1 小时数据源下线采集失败告警切换备用源2 小时4.2 数据延迟与不一致的处理金融数据对时效性要求高但数据源本身可能有延迟。比如某个数据源宣称实时行情延迟 3 秒实际高峰期可能延迟 30 秒。如果业务方按 3 秒的预期来用数据就会出问题。我的处理方式是在数据里带上数据时间戳和采集时间戳两个字段。数据时间戳是数据源标注的时间采集时间戳是我们实际拿到数据的时间。业务方可以根据这两个时间戳的差值判断数据新鲜度自行决定是否使用。对于多数据源的数据不一致问题我的策略是主源优先加交叉校验。每个数据类型指定一个主数据源其他源作为备份和校验。如果主源和备源的数据差异超过阈值记录异常并告警但默认仍然使用主源数据。这样既保证了服务可用性又能及时发现数据质量问题。注意不要试图自动合并多个数据源的数据除非你非常清楚每个源的质量特征。我试过用加权平均合并两个源的行情数据结果因为一个源在特定时段有系统性偏差合并后的数据反而更不准。后来改成主备模式简单可靠。4.3 数据库性能瓶颈的排查思路服务跑了一段时间后查询越来越慢这是必然会遇到的问题。排查数据库性能瓶颈我一般按这个顺序来先看慢查询日志再看索引使用情况最后看表膨胀和统计信息。慢查询日志是第一步。PostgreSQL 开启log_min_duration_statement后所有超过阈值的查询都会记录下来。我设的阈值是 200 毫秒大部分正常查询都在 50 毫秒以内超过 200 毫秒的基本都有优化空间。索引这块金融数据最常用的查询条件是“某个标的在某段时间范围内”。所以(symbol, trade_date)的复合索引是必须的而且trade_date要放在后面因为范围查询只能用在索引的最后一列。如果写成(trade_date, symbol)按标的查询时索引效果会差很多。表膨胀是 PostgreSQL 特有的问题。频繁的更新和删除会导致表和索引膨胀查询时需要扫描更多数据页。定期执行VACUUM ANALYZE能缓解这个问题但更根本的解决办法是分区。我按月分区后历史分区的数据基本不变膨胀问题自然就消失了。-- 查看表膨胀情况 SELECT schemaname, tablename, pg_size_pretty(pg_total_relation_size(schemaname||.||tablename)) AS total_size, pg_size_pretty(pg_relation_size(schemaname||.||tablename)) AS table_size FROM pg_tables WHERE schemaname public ORDER BY pg_total_relation_size(schemaname||.||tablename) DESC;4.4 服务高可用与故障恢复金融数据服务一旦挂掉依赖它的业务都会受影响。所以高可用设计是必须的。我的方案是多实例加健康检查加自动重启。多实例部署在至少两台机器上前面挂一个负载均衡。每个实例都暴露一个/health接口返回数据库连接状态、缓存连接状态、最近一次采集时间等信息。负载均衡定期检查这个接口发现异常就把流量切到其他实例。自动重启用 systemd 或 supervisor 都行我用的 systemd配置简单且和系统集成好。关键是重启策略要合理我设的是失败后 5 秒重启最多重启 5 次超过就停止并告警。不加限制的话如果服务因为配置错误反复崩溃会陷入无限重启循环。# /etc/systemd/system/financial-services.service [Unit] DescriptionFinancial Data Service Afternetwork.target postgresql.service redis.service [Service] Typeexec Userappuser WorkingDirectory/opt/financial-services ExecStart/opt/financial-services/venv/bin/gunicorn app.main:app \ --workers 4 \ --worker-class uvicorn.workers.UvicornWorker \ --bind 0.0.0.0:8000 Restarton-failure RestartSec5 StartLimitBurst5 [Install] WantedBymulti-user.target故障恢复方面最重要的是数据可重建。所有原始数据都有快照标准化数据可以从快照重新生成。所以即使数据库完全损坏也能在几个小时内恢复。我建议定期做恢复演练确保备份和恢复流程真的可用而不是等到出事才发现备份是坏的。4.5 常见问题速查表现象可能原因排查方法解决方案采集任务不执行调度器未启动查看调度器日志检查启动流程采集数据为空数据源接口变更对比原始快照更新适配器配置查询超时索引缺失或表膨胀查看慢查询日志加索引或执行 VACUUM缓存命中率低TTL 设置过短查看 Redis 统计调整 TTL 配置内存持续增长连接池泄漏查看连接数检查连接释放逻辑接口返回 500数据校验失败查看错误日志修复数据或放宽校验数据重复幂等逻辑缺失检查唯一约束加 ON CONFLICT 处理服务频繁重启内存不足查看系统日志增加内存或优化查询这套financial-services我从最初的一个脚本慢慢迭代到现在中间踩的坑基本都写在上面的内容里了。如果你刚开始搭类似的服务我的建议是先把适配器层和清洗层做扎实这两块做好了后面加数据源、改接口都会轻松很多。至于调度和缓存可以先用最简单的方案跑起来等遇到性能问题再优化不要一开始就过度设计。