ARTICLE DETAIL

资讯详情

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

从零搭建自托管金融数据服务:架构、采集、存储与接口设计

从零搭建自托管金融数据服务:架构、采集、存储与接口设计 1. 金融数据服务从零搭建的完整思路1.1 为什么我要自己搭一套金融数据服务先说清楚这个项目到底在干什么。financial-services这个名字听起来很泛实际上我做的是一套面向个人开发者和小型团队的自托管金融数据聚合与分发服务。核心功能就三件事从公开数据源定时抓取行情、财报、宏观经济指标把原始数据清洗成统一格式存进本地数据库对外暴露 REST 接口和 WebSocket 推送供自己的量化脚本、看板或者小工具消费。为什么不用现成的商业 API我算过一笔账。主流行情数据接口按调用次数计费一个中等频率的策略回测跑下来光数据成本就能吃掉大半利润。而且很多接口对历史数据的深度有限制分钟级数据往往只给最近几个月。自己搭一套前期投入大概两三天之后边际成本几乎为零数据想存多久存多久想怎么切怎么切。这套东西适合谁有一定编程基础、想认真做量化或者金融数据分析、又不愿意被数据费用卡脖子的个人开发者。如果你只是想看看大盘指数那直接用免费网页就够了没必要折腾这个。1.2 整体架构选型与背后的取舍架构上我走的是极简路线没有上微服务那一套。原因很直接个人项目的流量和并发量根本撑不起微服务的复杂度引入消息队列、服务发现这些组件只会让运维成本飙升。最终定下来的结构是四层采集层、存储层、计算层、服务层。采集层用 Python 写定时任务APScheduler做调度httpx做异步请求。选httpx而不是requests是因为它原生支持 async批量拉取几十个标的的行情时并发效率比同步请求高一个数量级。存储层用 PostgreSQL配合 TimescaleDB 扩展做时序数据。这里有个关键决策为什么不用 InfluxDB 或者 ClickHouseInfluxDB 在时序写入上确实快但它的 SQL 支持弱做关联查询比如把行情和财报数据 join 起来非常别扭。ClickHouse 查询性能强悍但对个人项目来说资源占用偏高单机跑起来内存吃紧。PostgreSQL TimescaleDB 的组合兼顾了时序写入效率和标准 SQL 的灵活性而且生态成熟遇到问题好查资料。计算层负责指标计算和数据对齐。不同数据源的时间戳精度不一样有的给到秒有的只给到天必须统一到同一时间轴上。服务层用 FastAPI自带 OpenAPI 文档省去写接口文档的功夫。WebSocket 推送用websockets库单独起一个轻量服务和 REST 接口解耦避免长连接把主服务的线程池占满。提示架构选型的第一原则是匹配自己的实际负载。个人项目上重型组件后期维护的时间成本远超收益。2. 数据采集环节的核心细节与实操要点2.1 数据源的选择与稳定性评估数据源是整个服务的命脉选错了后面全是坑。我评估数据源主要看四个维度数据覆盖面、更新频率、接口稳定性、使用条款的宽松程度。覆盖面决定了你能做什么策略更新频率决定了策略的时效性稳定性决定了你要不要写一堆重试逻辑使用条款则决定了你能不能合法地把数据用于自己的项目。我实际测试过七八个公开数据源最后保留了三类。第一类是官方交易所或监管机构提供的公开接口这类数据最权威但格式往往不统一需要写适配器。第二类是聚合型数据服务覆盖面广、格式统一但免费额度有限需要控制调用频率。第三类是财经资讯网站的公开页面作为补充数据源但要注意页面结构随时可能变解析逻辑要写得足够健壮。评估稳定性时我有个土办法连续跑一周的采集任务记录每次请求的成功率和响应时间画成曲线看波动。如果某个数据源的成功率低于 95%或者响应时间波动超过三倍我就会把它降级为备用源。这个测试成本很低但能避免上线后才发现数据源不可靠的尴尬。2.2 采集任务的调度设计与频率控制调度设计最容易犯的错误是把所有任务塞进一个定时器里。我一开始就是这么干的结果每次采集高峰期数据库连接池直接被打满其他任务全部超时。后来改成按数据源分组调度每组独立控制并发数问题就解决了。具体做法是用APScheduler的BlockingScheduler为每个数据源创建一个独立的 job每个 job 内部用asyncio.Semaphore限制并发请求数。比如行情数据更新频繁我设置每 5 分钟跑一次并发数限制在 10财报数据更新慢每天凌晨跑一次并发数限制在 3。这样即使某个数据源响应变慢也不会拖垮整个系统。频率控制还有个细节很多数据源对请求频率有隐性限制超过就会返回 429 或者直接封 IP。我的做法是在采集器里内置一个令牌桶限流器每个数据源配置独立的速率参数。这个参数不是拍脑袋定的而是根据数据源文档的说明再留出 30% 的余量。比如文档说每分钟最多 60 次请求我就设置成每分钟 40 次。实测下来这个余量能有效避免触发限流。import asyncio from aiolimiter import AsyncLimiter class DataCollector: def __init__(self, rate_limit: int, time_period: float 60.0): self.limiter AsyncLimiter(rate_limit, time_period) self.semaphore asyncio.Semaphore(10) async def fetch(self, url: str): async with self.limiter: async with self.semaphore: # 实际的请求逻辑 pass2.3 数据清洗与格式统一的实操方法原始数据拿到手只是第一步清洗才是真正花时间的活。不同数据源返回的字段名、时间格式、数值单位都不一样。比如有的用timestamp表示时间有的用date有的时间戳是秒级有的是毫秒级有的价格是字符串有的是浮点数。如果不统一后面查询和计算会痛苦不堪。我的清洗流程分三步。第一步是字段映射为每个数据源写一个映射配置把原始字段名映射到统一的内部字段名。这个配置用 YAML 文件管理改起来不用动代码。第二步是类型转换所有时间字段统一转成 UTC 时区的datetime对象所有数值字段统一转成Decimal类型。为什么用Decimal而不是float因为金融计算对精度极其敏感float的浮点误差在累加计算时会放大Decimal能保证精确。第三步是异常值处理比如价格出现负数或者单日涨跌幅超过 50%这些数据要么是错误要么是特殊事件我会打上标记存起来但不参与后续计算。注意清洗逻辑一定要写单元测试。我踩过的坑是某个数据源悄悄改了字段格式清洗代码没报错但数据全错了直到一周后才发现。现在每个数据源的清洗函数都有对应的测试用例每天采集完自动跑一遍。3. 存储层设计与数据模型落地3.1 时序数据表结构设计的关键决策存储层的核心是表结构设计。金融数据本质上是时序数据但又不是纯粹的时序数据因为它还涉及标的的元信息、财报的关联关系等。我最终设计了三类表标的元信息表、时序数据表、事件数据表。标的元信息表存股票代码、名称、所属行业、上市日期这些不常变的信息。时序数据表存行情、指标这类按时间排列的数据。事件数据表存分红、拆股、财报发布这类离散事件。为什么要把事件数据单独拆出来因为事件的查询模式和时序数据完全不同混在一起会导致索引效率下降。时序数据表用 TimescaleDB 的 hypertable 特性按时间自动分区。分区间隔我设置成一个月这个粒度是权衡的结果太细会导致分区数量爆炸太粗则查询时扫描的数据量太大。对于个人项目的数据量级按月分区是比较舒服的选择。主键用(symbol, timestamp)复合主键这样按标的和时间范围查询时能直接命中索引。CREATE TABLE market_data ( symbol VARCHAR(20) NOT NULL, timestamp TIMESTAMPTZ NOT NULL, open DECIMAL(18, 6), high DECIMAL(18, 6), low DECIMAL(18, 6), close DECIMAL(18, 6), volume BIGINT, PRIMARY KEY (symbol, timestamp) ); SELECT create_hypertable(market_data, timestamp, chunk_time_interval INTERVAL 1 month);3.2 数据写入的批量优化与去重策略写入性能是存储层的另一个关键点。逐条插入在数据量小的时候没问题但当天数据积累到百万行级别逐条插入会慢到无法接受。我的做法是批量插入每批 1000 行用execute_values或者COPY命令。实测下来批量插入比逐条插入快 20 倍以上。去重是必须处理的。采集任务可能因为重试或者调度重叠导致重复数据。我在表上建了唯一约束插入时用ON CONFLICT DO NOTHING或者ON CONFLICT DO UPDATE。前者用于行情数据因为同一时间点的数据应该是一样的后者用于财报数据因为财报可能修正需要更新为最新值。还有个细节是写入时间戳的处理。我额外加了一个created_at字段记录数据入库时间和业务时间戳分开。这样排查问题时能清楚知道数据是什么时候进来的而不是只知道数据对应的时间点。这个字段在调试采集延迟问题时特别有用。3.3 数据保留与归档的自动化方案数据不能无限存下去尤其是分钟级数据一年下来就是几千万行。我设计了一套分级保留策略分钟级数据保留最近 3 个月小时级数据保留最近 2 年日级数据永久保留。这个策略是根据实际使用场景定的——回测高频策略用最近几个月的数据就够了长期分析用日级数据。归档用 TimescaleDB 的原生压缩功能把 3 个月前的分钟级数据压缩存储压缩率大概能到 90% 以上。压缩后的数据仍然可以查询只是写入会被禁止。自动化用定时任务实现每周日凌晨跑一次归档脚本把符合条件的数据块压缩掉。这个操作对在线查询几乎没有影响因为 TimescaleDB 的压缩是在 chunk 级别进行的。4. 服务层接口设计与性能调优4.1 REST 接口的路径规划与参数设计服务层是这套系统对外的门面接口设计得好不好直接决定了用起来顺不顺手。我遵循的原则是路径表达资源参数表达过滤条件。比如获取行情数据的接口是GET /api/v1/market/{symbol}查询参数用start、end、interval来控制时间范围和粒度。参数设计上有个容易忽略的点默认值的选择。start默认值我设成当天零点end默认值设成当前时间interval默认值设成1d。这样即使调用方什么参数都不传也能拿到一份合理的默认数据降低了使用门槛。另外所有时间参数都接受 ISO 8601 格式也接受 Unix 时间戳内部统一转换处理。分页是必须的。行情数据动辄几千条一次性返回会拖慢响应。我用limit和offset做分页默认limit是 500最大允许 5000。超过最大值的请求会被拒绝并返回明确的错误信息而不是默默截断。这样调用方能清楚知道自己的请求是否被完整处理。4.2 查询性能优化的具体手段查询性能优化我做了三件事。第一是索引优化除了主键索引还在symbol和timestamp上分别建了索引因为查询模式既有按标的查全部历史也有按时间查所有标的。第二是查询缓存对于不常变的数据比如日级行情用 Redis 缓存查询结果缓存有效期设成 1 小时。第三是慢查询监控所有执行时间超过 500ms 的查询都会被记录到日志定期 review 并优化。这里重点说下缓存策略。缓存 key 的设计很关键我用md:{symbol}:{interval}:{start}:{end}作为 key这样不同参数组合的查询结果互不干扰。缓存失效用主动失效加被动过期结合的方式数据更新时主动删除相关 key同时设置过期时间兜底。实测下来加了缓存之后重复查询的响应时间从 200ms 降到了 5ms 以内。4.3 WebSocket 实时推送的实现细节实时推送是这套服务的亮点功能。实现上用websockets库起一个独立服务客户端连接后可以订阅特定标的的行情更新。推送逻辑是采集层写入新数据后通过 Redis 的 pub/sub 发一条消息WebSocket 服务订阅这个消息并推送给对应的客户端。连接管理上有个坑要注意客户端可能因为网络问题断开但服务端不知道导致连接泄漏。我的做法是加心跳机制服务端每 30 秒发一次 ping客户端必须回 pong连续三次没回应就主动断开。同时限制单个 IP 的最大连接数防止恶意连接耗尽资源。这些细节在文档里通常不会写但不做的话服务跑几天就会出问题。5. 常见问题排查与避坑经验实录5.1 数据采集失败的排查思路采集失败是最常见的问题排查时我按这个顺序走先看网络连通性再看数据源是否改了接口最后看自己的代码逻辑。网络问题好判断直接 curl 一下目标地址就知道。接口变更比较隐蔽通常表现为返回 200 但数据结构变了或者返回 404。我的做法是每次采集后校验数据条数和字段完整性异常时发告警。有个坑我踩过两次数据源在特定时间段比如收盘后会返回空数据但 HTTP 状态码是 200。如果代码只判断状态码就会把空数据当成正常数据存进去导致后续计算出现缺口。现在的做法是加一层数据质量检查如果某次采集的数据量比历史平均值低 50% 以上就标记为可疑并告警。5.2 数据库连接池耗尽的解决方案连接池耗尽通常发生在采集高峰期。表现是请求全部卡住日志里全是连接超时。根本原因是并发任务数超过了连接池大小。解决方案有两个方向一是调大连接池二是控制并发数。我倾向于后者因为连接池不是越大越好PostgreSQL 每个连接都有内存开销连接太多反而会拖慢整体性能。具体做法是给每个采集任务组分配独立的连接池池大小根据任务的实际并发需求设置。比如行情采集并发高给 20 个连接财报采集并发低给 5 个连接。同时设置连接的最大存活时间避免长时间空闲的连接占用资源。这个调整之后连接池耗尽的问题再没出现过。5.3 时间时区处理引发的数据错乱时区问题是金融数据里最隐蔽的坑。不同数据源用的时区不一样有的用 UTC有的用交易所所在地时区有的用北京时间。如果不统一跨数据源关联时就会出现时间对不上的情况。我的原则是存储层一律用 UTC展示层再根据用户需求转换。转换过程中有个细节夏令时。有些市场有夏令时切换切换当天的时间处理特别容易出错。我的做法是用pytz或者zoneinfo库处理时区转换不要自己手动加减小时数。另外所有涉及时间的比较和计算都先转成 UTC 再操作避免混用不同时区的时间对象。常见问题典型表现排查方向解决方案采集失败数据缺失、告警触发网络、接口变更、代码逻辑加数据质量校验异常告警连接池耗尽请求卡住、连接超时并发数超过池大小分组连接池控制并发时区错乱跨源数据时间对不上时区未统一存储统一 UTC用库转换写入缓慢采集任务积压逐条插入、索引过多批量插入精简索引缓存不一致查询结果过期失效策略不完善主动失效加被动过期5.4 服务上线的检查清单与个人体会上线前我一定会过一遍检查清单数据库索引是否齐全、连接池配置是否合理、限流参数是否设置、告警是否配置、备份是否正常。这份清单是踩坑踩出来的每一条背后都有一次事故。比如有一次忘了配告警采集任务挂了三天才发现数据缺口补起来非常麻烦。最后分享一个个人体会这套系统最大的价值不在于技术多先进而在于它完全受我控制。数据怎么存、接口怎么设计、什么时候更新全由自己决定。这种掌控感是使用商业服务给不了的。后续我打算加一个简单的回测框架直接消费这套服务的数据把从数据到策略的链路彻底打通。如果你也在做类似的事情建议先把采集和存储做扎实这两块稳了上层应用怎么折腾都不会出大问题。
返回列表