ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

从零搭建金融数据服务:架构设计与核心模块实现

从零搭建金融数据服务:架构设计与核心模块实现 1. 金融数据服务从零搭建的完整思路1.1 为什么我要自己搭一套金融数据服务先说清楚这个项目到底在干什么。financial-services这个名字听起来很泛但落到实际工程里它指的是一套面向金融场景的数据服务层——把行情、财报、汇率、宏观指标这些散落在各处的数据统一采集、清洗、存储再通过一套稳定的接口对外提供查询能力。它解决的核心问题是当你需要在自己的应用里接入金融数据时不用每次都去对接不同的数据源、处理不同的格式、重复写解析逻辑而是有一个统一的入口。这套东西适合谁如果你在做量化回测、做投资组合看板、做企业财务分析工具或者只是想在个人项目里加一个实时汇率或股票行情的小模块那这套服务的思路都能直接复用。哪怕你只是刚入门的开发者只要会写基本的增删改查跟着走一遍也能搭出一个能跑的最小版本。我之所以决定自己搭而不是直接用现成的商业接口原因很实际一是成本很多数据源按调用次数收费量一大就吃不消二是可控性商业接口说改就改、说停就停数据格式一变你的代码就得跟着改三是学习价值自己走一遍采集到服务的全链路对金融数据的理解会完全不一样。当然这不是说商业接口不好而是说在特定场景下自建服务有它不可替代的优势。1.2 整体架构怎么设计才不返工搭这类服务最容易犯的错就是一上来就写代码写到一半发现数据结构不对、扩展不了只能推倒重来。我的经验是先把架构想清楚哪怕只是在纸上画一画。整体上我把它分成四层采集层、存储层、服务层、调度层。采集层负责从各个数据源拉数据存储层负责落地和缓存服务层对外提供接口调度层负责定时任务和失败重试。这四层之间通过明确的接口通信任何一层换实现都不影响其他层。为什么这么分因为金融数据有个特点——数据源极不稳定。今天能用的接口明天可能就限流了这个源的数据格式下个月可能就变了。如果采集逻辑和业务逻辑混在一起每次数据源变动都要动核心代码维护成本会爆炸。分层之后采集层可以随时替换某个源只要输出格式统一上层完全无感。存储层我选的是PostgreSQL Redis的组合。PostgreSQL 存历史数据和结构化数据比如日线行情、财报数据Redis 做热点缓存比如最新报价、常用查询结果。为什么不用纯 Redis因为金融数据需要复杂查询和事务保证Redis 在这块能力有限。为什么不用纯 PostgreSQL因为高频查询场景下每次都打数据库扛不住缓存能挡掉大部分重复请求。服务层用FastAPI搭原因是它自带异步支持、自动生成接口文档、类型校验完善对于数据服务这种 IO 密集型的场景非常合适。调度层用APScheduler轻量、够用不需要引入 Celery 那种重量级方案。提示架构设计阶段不要追求完美而是追求可替换。每个模块都留好替换的口子后面改起来才不痛苦。1.3 数据源选型的取舍逻辑金融数据源大致分三类公开免费源、商业付费源、自建采集源。公开免费源比如一些开放数据平台优点是零成本缺点是稳定性差、频率限制严、字段不全。商业付费源质量高但贵。自建采集源灵活但维护成本高。我的策略是混合使用核心数据用一两个稳定的源做主力备用源做兜底非核心数据用免费源。具体到实现上每个数据源都封装成一个独立的适配器统一实现fetch()和parse()两个方法。这样新增一个源只需要写一个适配器不用动其他代码。选源的时候有几个硬指标必须看更新频率、历史深度、字段完整度、调用限制、数据格式稳定性。我踩过的坑是只看更新频率结果用了一个历史数据只有半年的源做回测的时候发现数据不够只能临时换源浪费了两天。所以选源之前一定要把这几项都确认清楚最好先跑一周的测试数据看看稳定性。2. 核心模块拆解与关键实现细节2.1 采集层怎么把脏数据变成干净数据采集层是整个服务的地基也是最脏最累的活。金融数据的脏体现在几个方面格式不统一、字段缺失、时间戳时区混乱、数值精度不一致。比如同样是收盘价有的源给的是字符串有的是浮点数有的还带货币符号。我的处理流程是三步走拉取、标准化、校验。拉取阶段只负责把原始数据拿回来不做任何处理原样存一份到临时表。标准化阶段把数据转成统一的内部格式包括时间统一转 UTC、数值统一转 Decimal、字段名统一映射。校验阶段检查数据合理性比如价格不能为负、时间不能是未来、涨跌幅不能超过合理范围。为什么用 Decimal 而不是 float因为金融计算对精度极其敏感。float 在二进制下无法精确表示很多十进制小数累加之后误差会放大。Decimal 虽然慢一点但精度可控做金额计算必须用它。这个细节很多人会忽略等到对账的时候发现差几分钱排查起来非常痛苦。from decimal import Decimal, ROUND_HALF_UP def normalize_price(raw_value): # 去掉货币符号和千分位 cleaned str(raw_value).replace(,, ).replace($, ).strip() # 统一保留4位小数四舍五入 return Decimal(cleaned).quantize(Decimal(0.0001), roundingROUND_HALF_UP)校验这块我建议宁可拒绝也不要存疑。遇到不合理的数据直接丢弃并记录日志不要试图修复它。因为金融数据一旦错了下游所有分析都是错的而且很难追溯。我见过有人为了不浪费数据把异常值强行修正结果回测结果完全失真白忙一场。2.2 存储层表结构设计的几个关键决策表结构设计直接决定了后续查询的性能和扩展性。金融数据有几个特点写入频繁、查询模式固定、历史数据量大。针对这些特点我的设计原则是读写分离、冷热分离、按时间分区。以行情数据为例我建了两张表quotes_realtime存最新报价只保留每个标的一条记录用 Redis 做缓存quotes_history存历史数据按日期分区每个月一个分区。为什么这么分因为实时查询和历史查询的模式完全不同——实时查询要快历史查询要全。混在一张表里两边都做不好。分区的好处是查询时能自动裁剪只扫描相关分区速度提升非常明显。我实测过一张 5000 万行的行情表不分区查询要十几秒按天分区后同样的查询只要几百毫秒。分区键选时间是因为金融查询几乎都带时间范围这是最自然的裁剪维度。索引方面我建了标的代码 时间的联合索引覆盖大部分查询场景。注意索引不是越多越好每个索引都会拖慢写入速度。金融数据写入量大索引要精打细算。我的做法是先上线观察慢查询日志再针对性加索引而不是一开始就堆一堆。表名用途分区策略索引quotes_realtime最新报价不分区标的代码唯一索引quotes_history历史行情按月分区标的代码时间联合索引fundamentals财报数据按年分区标的代码报告期索引macro_indicators宏观指标不分区指标代码时间索引2.3 服务层接口设计怎么兼顾灵活和性能服务层的接口设计要回答一个问题调用方到底需要什么。我观察下来金融数据查询无非几种模式按标的查最新、按标的查历史区间、按条件批量查、按指标聚合查。针对这几种模式设计接口比设计一个万能接口要好得多。万能接口的问题是参数太多、语义模糊、难以优化。比如一个接口同时支持查股票和查汇率内部就要写一堆 if-else缓存也不好做。我倾向于按数据域拆分接口行情一个、财报一个、宏观一个每个接口内部逻辑清晰缓存策略也能针对性设计。性能优化上我做了三件事缓存、批量、异步。缓存用 Redis热点数据 TTL 设短一点比如 5 秒冷数据设长一点比如 1 小时。批量是指接口支持一次查多个标的减少往返次数。异步是指用 FastAPI 的 async 特性IO 等待时不阻塞其他请求。from fastapi import FastAPI, Query from typing import List app FastAPI() app.get(/quotes/latest) async def get_latest_quotes(symbols: List[str] Query(...)): # 先查缓存未命中的再查数据库 result {} missing [] for sym in symbols: cached await redis.get(fquote:{sym}) if cached: result[sym] cached else: missing.append(sym) if missing: db_data await fetch_from_db(missing) result.update(db_data) # 回写缓存 for sym, val in db_data.items(): await redis.setex(fquote:{sym}, 5, val) return result注意缓存回写一定要设 TTL否则数据更新后缓存永远不失效会返回过期数据。这个坑我踩过排查了半天才发现是缓存没设过期时间。2.4 调度层定时任务怎么跑才不出乱子调度层负责定时拉数据、重试失败任务、清理过期数据。看起来简单但实际跑起来问题不少。最常见的是任务堆积——上一个任务还没跑完下一个又开始了最后把系统拖垮。我的解决方案是加锁 超时。每个任务执行前先抢一个分布式锁抢不到就跳过本次执行。任务设置最大执行时间超时强制终止并记录。这样即使某个任务卡住也不会影响后续任务。锁用 Redis 实现简单可靠。重试策略上我用的是指数退避。第一次失败等 1 秒重试第二次等 2 秒第三次等 4 秒最多重试 5 次。为什么不用固定间隔因为很多失败是临时的比如网络抖动、对方限流固定间隔重试容易撞上对方的限流窗口指数退避能错开。重试次数也不能太多否则一个坏任务会一直占着资源。import asyncio from redis import Redis redis_client Redis() async def run_with_lock(task_name, task_func, timeout300): lock_key flock:{task_name} # 抢锁过期时间设长一点防止死锁 if not redis_client.set(lock_key, 1, nxTrue, extimeout): return # 已有任务在跑跳过 try: await asyncio.wait_for(task_func(), timeouttimeout) except asyncio.TimeoutError: log.error(f{task_name} 执行超时) finally: redis_client.delete(lock_key)3. 完整实操流程与关键环节落地3.1 环境准备与依赖安装动手之前先把环境理清楚。我用的是 Python 3.10因为要用到一些较新的语法特性。数据库 PostgreSQL 14Redis 6。这些版本不是随便定的——PostgreSQL 14 之后分区表的性能有明显提升Redis 6 之后支持多线程 IO对高并发场景更友好。依赖管理我用的是poetry比pip加requirements.txt更清晰能锁定依赖版本避免在我机器上能跑的问题。核心依赖就几个fastapi、uvicorn、asyncpg、redis、apscheduler、httpx、pydantic。httpx用来做异步 HTTP 请求比requests更适合异步场景。# 初始化项目 poetry new financial-services cd financial-services # 添加依赖 poetry add fastapi uvicorn asyncpg redis apscheduler httpx pydantic # 启动开发服务器 poetry run uvicorn app.main:app --reload数据库初始化我写了一个init.sql包含建表、建分区、建索引的语句。分区表用 PostgreSQL 的声明式分区每个月自动创建一个新分区。这里有个细节分区不会自动创建需要提前建好或者用定时任务建。我选择提前建一年的分区然后每月检查一次避免临时建分区时锁表。提示建分区的时候记得给分区也建索引否则查询时分区裁剪了但索引没生效性能还是上不去。这个坑很隐蔽因为主表的索引不会自动继承到分区。3.2 数据采集的完整实现采集流程我拆成四个步骤配置数据源、拉取原始数据、标准化处理、入库。配置数据源用一个 YAML 文件管理每个源配 URL、频率限制、字段映射。这样新增源不用改代码改配置就行。拉取环节要注意频率控制。很多免费源有调用频率限制超了会被封。我用令牌桶算法做限流每个源独立配置速率。令牌桶的好处是允许突发流量同时保证长期速率不超标。实现上用 Redis 的原子操作多实例部署时也能共享限流状态。标准化环节是重头戏。不同源的数据格式差异很大我定义了一套内部标准格式所有源的数据都往这个格式上靠。字段映射用配置驱动比如源 A 的close_price映射到内部的close源 B 的last也映射到close。这样上层完全不用关心数据来自哪个源。import yaml from decimal import Decimal class DataSourceAdapter: def __init__(self, config_path): with open(config_path) as f: self.config yaml.safe_load(f) async def fetch(self, symbol): url self.config[url_template].format(symbolsymbol) async with httpx.AsyncClient() as client: resp await client.get(url, timeout10) resp.raise_for_status() return resp.json() def parse(self, raw_data): mapping self.config[field_mapping] result {} for internal_field, source_field in mapping.items(): value raw_data.get(source_field) if internal_field in (open, high, low, close): value Decimal(str(value)).quantize(Decimal(0.0001)) result[internal_field] value return result入库环节用批量插入 冲突更新。金融数据经常需要更新比如盘中价格变动用INSERT ... ON CONFLICT DO UPDATE一条语句搞定比先查后插效率高得多。批量插入时每批 1000 条左右太大容易超时太小效率低。3.3 服务接口的完整实现接口实现我遵循薄接口、厚服务的原则。接口层只做参数校验和响应封装业务逻辑放在服务层。这样接口层可以随时替换比如从 REST 换成 GraphQL业务逻辑不用动。以查询历史行情为例接口接收标的代码、开始时间、结束时间、频率四个参数。频率参数决定返回日线、周线还是月线。服务层根据频率决定是直接查表还是做聚合。日线直接查周线和月线用 PostgreSQL 的date_trunc聚合。app.get(/quotes/history) async def get_history( symbol: str, start: str, end: str, freq: str daily ): # 参数校验 if freq not in (daily, weekly, monthly): raise HTTPException(400, freq 参数不合法) # 根据频率选择查询策略 if freq daily: rows await db.fetch( SELECT date, open, high, low, close FROM quotes_history WHERE symbol $1 AND date BETWEEN $2 AND $3 ORDER BY date, symbol, start, end ) else: trunc week if freq weekly else month rows await db.fetch( fSELECT date_trunc({trunc}, date) as period, first(open, date) as open, max(high) as high, min(low) as low, last(close, date) as close FROM quotes_history WHERE symbol $1 AND date BETWEEN $2 AND $3 GROUP BY period ORDER BY period, symbol, start, end ) return {symbol: symbol, freq: freq, data: rows}响应格式我统一用{code, message, data}的结构。code是业务状态码message是提示信息data是实际数据。这样调用方处理起来一致不用为每个接口写不同的解析逻辑。错误码我定义了一套比如 1001 表示参数错误2001 表示数据源不可用3001 表示内部错误。3.4 调度任务的完整配置调度任务用 APScheduler 配置我定义了三个核心任务行情采集每分钟、财报采集每天、数据清理每周。频率不同用不同的调度器实例避免相互影响。行情采集每分钟跑一次因为盘中价格变动频繁。财报采集每天凌晨跑一次因为财报更新频率低。数据清理每周跑一次删除过期缓存和临时数据。每个任务都配了重试和告警失败时发通知。from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger scheduler AsyncIOScheduler() # 行情采集交易时段每分钟 scheduler.add_job( collect_quotes, CronTrigger(day_of_weekmon-fri, hour9-15, minute*), idcollect_quotes, max_instances1, # 防止任务堆积 misfire_grace_time30 # 错过执行的宽限时间 ) # 财报采集每天凌晨2点 scheduler.add_job( collect_fundamentals, CronTrigger(hour2, minute0), idcollect_fundamentals, max_instances1 ) # 数据清理每周日凌晨3点 scheduler.add_job( cleanup_old_data, CronTrigger(day_of_weeksun, hour3, minute0), idcleanup ) scheduler.start()max_instances1这个参数很关键它保证同一个任务不会并发执行。misfire_grace_time是错过执行的宽限时间比如服务器重启导致任务错过30 秒内还能补跑。这两个参数不设的话任务堆积和重复执行的问题迟早会出现。4. 常见问题排查与避坑经验实录4.1 数据源相关的典型问题数据源问题占了实际运维问题的一大半。最常见的是接口突然不可用表现为超时、返回 403、返回格式变化。我的排查顺序是先看是不是网络问题ping 一下再看是不是限流看返回头最后看是不是格式变了打印原始响应。限流问题最隐蔽因为很多源不返回明确的限流提示只是默默返回空数据或旧数据。我的应对是监控数据新鲜度——如果某个源的数据超过预期时间没更新就告警。这个监控比监控接口状态更有效因为接口可能返回 200 但数据是旧的。格式变化也很常见尤其是免费源。我的应对是严格校验 快速失败。解析时如果发现字段缺失或类型不对直接抛异常并记录原始数据不要试图猜测。这样问题能第一时间暴露而不是等到下游分析出错才发现。问题现象可能原因排查方法解决方案接口超时网络问题或对方限流ping 看返回头重试 降级到备用源返回 403频率超限或 IP 被封检查调用频率降低频率 换源数据不更新源停止更新或解析失败对比原始数据告警 人工介入字段缺失源格式变化打印原始响应更新字段映射配置数值异常源数据错误校验规则检查丢弃 记录日志4.2 性能问题的排查思路性能问题通常表现为接口变慢、数据库 CPU 高、缓存命中率低。我的排查顺序是先看监控、再看慢查询、最后看代码。监控能快速定位是哪个环节慢慢查询日志能定位具体 SQL代码审查能发现逻辑问题。数据库慢查询是最常见的性能瓶颈。金融数据的查询往往涉及大范围扫描如果不走索引就会很慢。我的经验是用 EXPLAIN ANALYZE 看执行计划确认索引是否生效。有时候索引建了但没用上是因为查询条件类型不匹配比如字符串和数字比较这种问题很隐蔽。缓存命中率低也要重视。命中率低意味着大量请求打到数据库数据库压力大。原因可能是 TTL 设太短、缓存键设计不合理、或者缓存被频繁失效。我的做法是监控命中率低于 80% 就排查。缓存键我统一用业务:标的:时间粒度的格式清晰且不易冲突。注意缓存雪崩是个大坑。如果大量缓存同时过期请求会瞬间全打到数据库。解决办法是给 TTL 加随机抖动比如 5 秒的 TTL 实际设成 4 到 6 秒之间的随机值避免同时过期。4.3 数据一致性的保障技巧金融数据对一致性要求高但分布式环境下一致性很难保证。我的策略是最终一致 对账。采集和入库不是原子的可能采到了但没入库或者入库了但采集标记没更新。这些不一致通过定时对账来发现和修复。对账的逻辑是对比源数据和本地数据找出差异并修复。对账频率不用太高每天一次就够。对账时要注意时间窗口因为源数据可能还在更新对账要选已经稳定的时间段。比如对账昨天的数据而不是今天的数据。另一个一致性问题是并发写入。多个采集任务同时写同一个标的的数据可能产生冲突。我的解决办法是按标的加锁同一个标的同一时间只有一个任务在写。锁的粒度要细按标的加锁而不是全局加锁否则并发度太低。async def safe_upsert(symbol, data): lock_key fwrite_lock:{symbol} # 尝试获取写锁最多等5秒 acquired await redis.set(lock_key, 1, nxTrue, ex5) if not acquired: log.warning(f{symbol} 正在被其他任务写入跳过) return False try: await db.execute( INSERT INTO quotes_history (symbol, date, close) VALUES ($1, $2, $3) ON CONFLICT (symbol, date) DO UPDATE SET close EXCLUDED.close, symbol, data[date], data[close] ) return True finally: await redis.delete(lock_key)4.4 独家避坑经验汇总最后分享几个我在实际项目中踩过的坑都是文档里不会写的。第一个坑时区问题。金融数据的时间戳时区五花八门有的用 UTC有的用交易所本地时间有的干脆不标时区。我一开始没注意导致不同源的数据时间对不上做聚合时出现重复或缺失。后来统一规定入库前全部转 UTC展示时再转本地。这个规则看似简单但执行起来要每个源都检查漏一个就出问题。第二个坑精度问题。前面提过用 Decimal但还有个细节——Decimal 的上下文精度。Python 的 Decimal 默认精度是 28 位做除法时可能不够。我遇到过计算收益率时精度丢失导致结果偏差。解决办法是显式设置上下文精度或者用quantize控制结果精度。第三个坑连接池耗尽。数据库连接池大小设小了高并发时请求排队设大了数据库扛不住。我的经验是从 10 开始根据监控调整。同时要确保连接用完及时归还异常时也要归还。用async with上下文管理器能自动处理比手动管理可靠。第四个坑日志太多。采集任务每分钟跑一次每次都打日志日志文件很快就爆了。我的做法是分级日志 采样。正常情况打 DEBUG 级别异常打 ERROR 级别。高频任务的日志做采样比如每 100 次打一次避免日志淹没真正的问题。第五个坑配置硬编码。一开始图省事把数据源 URL、频率限制这些写死在代码里。后来要改一个参数就得重新部署非常麻烦。后来全部抽到配置文件支持热加载改配置不用重启。这个改动虽然前期多花时间但后期省了无数事。这套服务我陆陆续续迭代了大半年从最初只能拉一个源的数据到现在支持多个源、多种数据、完整的监控告警。回头看最难的不是写代码而是想清楚数据怎么流动、异常怎么处理、扩展怎么做。代码只是把这些想清楚的东西表达出来。如果你也在搭类似的东西我的建议是先把数据流图画清楚再动手写代码能省掉大量返工。
返回列表