1. 金融数据服务从零搭建的核心思路
1.1 为什么选这个方向
金融数据服务这个方向,说白了就是解决一个很朴素的问题:数据从哪来、怎么存、怎么算、怎么给出去。我最早接触这块是因为帮一个做量化的小团队搭后台,他们每天要处理几十万条行情快照和基本面数据,一开始用Excel加脚本硬扛,后来数据量一上来直接崩了。这个痛点非常典型——数据源分散、格式不统一、查询慢、扩展难。
做金融数据服务,核心目标就三个:数据准确、响应快、能扛住并发。适合谁来参考?如果你正在做量化交易系统、金融风控平台、或者任何需要处理时序数据的项目,这套思路都能直接套用。哪怕你只是想把股票数据拉下来做个人分析,里面的分层设计和缓存策略也值得一看。
1.2 整体架构怎么分层
我习惯把金融数据服务拆成四层,每层职责单一,方便替换和扩展:
- 采集层:负责从各种数据源拉数据,包括交易所接口、公开API、文件导入等。这一层的关键是容错和去重,因为金融数据源经常抽风,重复推送是家常便饭。
- 存储层:时序数据用列式存储,关系型数据用传统数据库。我一般会做冷热分离,最近3个月的热数据放内存或SSD,历史数据归档到便宜的对象存储。
- 计算层:指标计算、聚合、回测引擎都在这层。重点是增量计算,别每次全量重算,否则数据量一大直接卡死。
- 服务层:对外提供RESTful API或WebSocket推送,做限流、鉴权、缓存。
这个分层的好处是,哪天你想换个数据源,只动采集层就行;想换存储引擎,也只影响存储层。解耦是金融系统能长期维护的前提,我见过太多项目因为采集和计算揉在一起,最后改一行代码牵一发动全身。
1.3 技术选型的取舍逻辑
选型这块我踩过不少坑,说几个关键决策点:
数据库选型:时序数据我首推ClickHouse,写入快、压缩率高、聚合查询性能炸裂。但如果你的团队对运维复杂度敏感,TimescaleDB(基于PostgreSQL)是更稳妥的选择,生态成熟,SQL兼容性好。关系型数据就用PostgreSQL,别用MySQL了,金融场景下PostgreSQL的窗口函数、JSON支持、并发控制都更胜一筹。
消息队列:Kafka是标配,但如果你数据量不大(每天百万级以下),Redis Stream或者RabbitMQ完全够用,别为了用Kafka而用Kafka,运维成本摆在那。
计算引擎:Python pandas适合中小规模,但数据量上到千万级就得换Polars或者直接上Spark。我实测下来,Polars在单机上的性能比pandas快5到10倍,内存占用还低,中小团队强烈推荐。
提示:选型时优先考虑团队熟悉的技术栈,而不是盲目追新。金融系统稳定压倒一切,一个你完全掌握的普通方案,比一个你半懂不懂的先进方案靠谱得多。
2. 数据采集与清洗的实操细节
2.1 数据源接入的常见坑
金融数据源大致分三类:交易所直连、第三方API、文件导入。交易所直连延迟最低但接入复杂,第三方API最方便但有速率限制和费用,文件导入适合历史数据补录。
我重点说第三方API接入的坑。第一,速率限制,很多免费API每分钟只给几十次请求,你得做令牌桶限流,别傻乎乎地循环调用,封IP是分分钟的事。第二,数据格式不统一,同一个字段不同源可能叫close、close_price、closingPrice,你得建一个字段映射表。第三,时间戳时区,这是最容易被忽略的,有的源给UTC,有的给北京时间,不统一处理后面计算全乱套。
# 字段映射表示例 FIELD_MAPPING = { "close": ["close", "close_price", "closingPrice", "收盘价"], "volume": ["volume", "vol", "trade_volume", "成交量"], "timestamp": ["timestamp", "time", "ts", "datetime"] } def normalize_record(raw, mapping): normalized = {} for std_field, aliases in mapping.items(): for alias in aliases: if alias in raw: normalized[std_field] = raw[alias] break return normalized2.2 数据清洗的五个关键步骤
采集来的原始数据基本不能直接用,我一般走这五步:
- 去重:按主键(通常是标的代码+时间戳)去重,保留最新一条。用数据库的
ON CONFLICT或者Redis的Set都能做。 - 缺失值处理:金融数据缺失很常见,停牌、网络抖动都会导致。我的原则是前向填充为主,插值为辅,但要在数据里标记哪些是填充的,别让下游误以为是真实数据。
- 异常值检测:价格突然涨跌超过阈值(比如10%)要标记出来人工复核,别自动删,万一是真实行情呢。
- 时间对齐:不同频率的数据要统一到同一时间轴,比如日线数据和分钟数据对齐,用
resample做重采样。 - 标准化:字段类型统一、单位统一(比如成交量统一成股还是手)、精度统一。
注意:清洗规则一定要版本化,每次改动都记录在案。我吃过亏,改了清洗逻辑后历史数据和新增数据口径不一致,排查了整整两天。
2.3 增量采集与断点续传
全量采集只适合初始化,日常必须走增量。核心是记录水位线(watermark),每次采集完更新水位线,下次从水位线之后开始拉。
断点续传的关键是幂等性,同一批数据重复写入不能产生副作用。我的做法是给每条记录算一个唯一哈希(标的+时间戳+关键字段),写入时用INSERT ... ON CONFLICT DO NOTHING,重复的直接跳过。
-- ClickHouse的幂等写入示例 INSERT INTO market_data SELECT * FROM staging_table WHERE (symbol, timestamp) NOT IN (SELECT symbol, timestamp FROM market_data);这套机制实测下来很稳,哪怕采集程序半夜崩了,重启后自动从断点继续,不会丢数据也不会重复。
3. 存储设计与查询优化实战
3.1 时序数据的表结构设计
金融时序数据的表结构设计直接决定查询性能。我推荐宽表+分区的方案:
CREATE TABLE market_data ( symbol String, trade_date Date, trade_time DateTime, open Float64, high Float64, low Float64, close Float64, volume UInt64, amount Float64, adj_factor Float64 ) ENGINE = MergeTree() PARTITION BY toYYYYMM(trade_date) ORDER BY (symbol, trade_time) SETTINGS index_granularity = 8192;几个关键点:分区键用月,别用天,否则分区太多元数据爆炸;排序键用symbol+time,因为查询基本都是按标的和时间范围来的;index_granularity默认8192就行,调太小索引膨胀,调太大扫描变慢。
3.2 查询优化的几个狠招
金融数据查询有两个典型场景:单标的时序查询和多标的截面查询。优化手段不一样。
单标的时序查询,靠排序键就能搞定,ClickHouse的稀疏索引直接定位到数据块。多标的截面查询(比如查某天所有股票的收盘价),需要预聚合或者物化视图。
-- 物化视图:按日预聚合 CREATE MATERIALIZED VIEW daily_summary ENGINE = SummingMergeTree() PARTITION BY toYYYYMM(trade_date) ORDER BY (trade_date, symbol) AS SELECT trade_date, symbol, argMax(close, trade_time) AS close, sum(volume) AS volume, sum(amount) AS amount FROM market_data GROUP BY trade_date, symbol;物化视图的代价是写入放大,但查询性能提升是数量级的。我实测过一个场景,原来查全市场某天的数据要3秒,加了物化视图后50毫秒。
3.3 冷热分离与数据归档
金融数据的特点是越老的数据查得越少。我一般做三级存储:
| 数据年龄 | 存储介质 | 查询延迟 | 成本 |
|---|---|---|---|
| 0-3个月 | SSD/内存 | 毫秒级 | 高 |
| 3个月-2年 | 普通磁盘 | 秒级 | 中 |
| 2年以上 | 对象存储 | 分钟级 | 低 |
归档策略用定时任务,每月把超过3个月的数据从热存储迁到冷存储。查询时如果热存储没有,自动回源到冷存储。这套方案帮一个客户把存储成本降了60%,查询性能几乎没影响。
提示:归档前一定要做数据校验,我见过归档过程中因为编码问题导致数据损坏的案例,血的教训。
4. 服务层API设计与性能保障
4.1 API设计的三个原则
金融数据服务的API设计,我坚持三个原则:语义清晰、版本可控、限流明确。
语义清晰指的是URL和参数一看就懂,比如/api/v1/market/kline?symbol=000001&start=2024-01-01&end=2024-03-01&freq=1d,别搞什么/api/getData?type=1&p=xxx。
版本可控是指API必须带版本号,/api/v1/、/api/v2/,老版本至少保留半年,给下游迁移时间。
限流明确是指每个API都要有明确的QPS限制,并在响应头里返回剩余额度,让调用方心里有数。
4.2 缓存策略的分层设计
金融数据查询,缓存是性能的生命线。我一般做三层缓存:
- 本地缓存:用进程内的LRU缓存,存最近查询的热点数据,命中率能到40%左右。
- 分布式缓存:Redis存全量热点数据,TTL设短一点(比如5分钟),保证数据新鲜度。
- 数据库缓存:ClickHouse本身的查询缓存,对重复查询有效。
缓存更新的策略是写时失效+读时重建。数据更新时删掉对应缓存key,下次查询时重新从数据库加载。别用定时刷新,金融数据时效性要求高,定时刷新会导致数据不一致。
import redis import json from functools import lru_cache r = redis.Redis(host='localhost', port=6379, db=0) @lru_cache(maxsize=1000) def get_kline_cached(symbol, start, end, freq): cache_key = f"kline:{symbol}:{start}:{end}:{freq}" cached = r.get(cache_key) if cached: return json.loads(cached) data = query_from_db(symbol, start, end, freq) r.setex(cache_key, 300, json.dumps(data)) return data4.3 并发压力下的稳定性保障
金融数据服务经常面临突发流量,比如开盘瞬间、财报发布时。保障稳定性靠这几招:
限流:用令牌桶算法,每个用户分配独立的桶,防止单个用户打满。Nginx的limit_req或者应用层的ratelimit库都能做。
熔断:当数据库响应时间超过阈值,自动切断请求,返回缓存数据或降级结果。用Hystrix或者Resilience4j。
异步化:非实时查询走异步任务队列,用户提交任务后拿task_id轮询结果,别让HTTP连接一直挂着。
压测:上线前必须压测,我用Locust模拟过1000并发,发现连接池配置太小导致大量超时,调大连接池后QPS从200提升到1500。
注意:压测环境要和生产环境配置一致,我见过在开发机上压测通过,上线后直接崩的案例,因为开发机没开连接池限制。
5. 常见问题与排查技巧实录
5.1 数据不一致的排查思路
数据不一致是金融系统最头疼的问题,表现是同一指标不同地方查出来不一样。排查思路按这个顺序来:
- 确认时间范围:先看查询的时间范围是否一致,时区是否统一。
- 确认数据版本:有没有用到缓存,缓存是否过期。
- 确认计算逻辑:复权因子、汇率转换这些是否一致。
- 确认数据源:是不是从不同源查的,不同源数据本身就有差异。
我遇到过一次,两个系统查同一只股票的收盘价差0.01,最后发现是复权处理不同,一个用前复权一个用后复权。这种问题只能靠统一数据口径来解决,建一个数据字典,所有系统都按这个来。
5.2 性能突然下降的应急处理
性能突然下降,先别急着改代码,按这个顺序排查:
| 排查项 | 检查方法 | 常见原因 |
|---|---|---|
| 数据库连接 | 查连接池状态 | 连接泄漏、连接池太小 |
| 慢查询 | 开慢查询日志 | 缺索引、全表扫描 |
| 缓存命中率 | 看Redis监控 | 缓存穿透、key设计不合理 |
| 系统资源 | top/free/iostat | CPU打满、内存不足、磁盘IO瓶颈 |
| 网络 | ping/traceroute | 网络抖动、带宽打满 |
我处理过一次线上故障,查询延迟从50ms飙到5秒,最后定位是Redis内存满了触发淘汰,大量请求穿透到数据库。解决办法是加内存+优化key的TTL策略。
5.3 数据采集断流的恢复流程
采集断流的原因很多:API限额、网络中断、程序崩溃。恢复流程我总结成四步:
- 确认断流时间点:从监控告警或者日志里找到最后一次成功采集的时间。
- 检查数据源状态:确认是源的问题还是自己的问题。
- 补采数据:从断流时间点开始重新采集,注意去重。
- 校验数据完整性:补采后做一次全量校验,确认没有缺口。
补采的时候要注意别把补采和实时采集混在一起,我一般用独立的补采任务,补采完成后合并到主表。
提示:采集程序一定要加心跳监控,超过5分钟没数据就告警,别等下游发现数据不对才来查。
5.4 常见问题速查表
| 问题现象 | 可能原因 | 解决方法 |
|---|---|---|
| 查询超时 | 缺索引、数据量太大 | 加索引、加物化视图、分页查询 |
| 数据重复 | 采集幂等没做好 | 加唯一约束、用ON CONFLICT |
| 内存溢出 | 一次性加载太多数据 | 分批查询、用流式处理 |
| 写入慢 | 分区太多、索引太多 | 调整分区粒度、减少索引 |
| 缓存不一致 | TTL太长、更新策略不对 | 缩短TTL、写时失效 |
| API被刷 | 没限流 | 加令牌桶、加鉴权 |
这套速查表是我从多次故障中总结出来的,基本覆盖了80%的常见问题。遇到新问题先查表,查不到再深入排查。
6. 个人实操体会与扩展方向
6.1 几个让我印象深刻的坑
第一个坑是浮点数精度。金融数据用Float64存储,计算的时候会出现0.1+0.2=0.30000000000000004这种问题。后来我把价格字段改成Decimal类型,或者用整数存储(价格乘以10000存整数),彻底解决。
第二个坑是时区处理。早期没统一时区,导致跨市场数据对齐时差了好几个小时。后来所有时间戳统一存UTC,展示的时候再转本地时区。
第三个坑是批量写入的事务问题。一次写入10万条数据,中途失败导致部分写入,数据不一致。后来改成小批量写入,每批1000条,失败只影响当前批次。
6.2 后续可以扩展的方向
这套架构搭好后,可以往几个方向扩展:
实时计算:接入Flink或Spark Streaming,做实时指标计算和告警。比如价格突破阈值实时推送。
机器学习:在计算层加模型推理,做价格预测、异常检测。特征工程可以直接复用存储层的数据。
多市场支持:扩展到期货、外汇、加密货币,核心架构不用变,只需要加采集适配器和字段映射。
数据可视化:对接Grafana或自研前端,做实时监控和交互式分析。
我个人在实际操作中的体会是,金融数据服务最难的不是技术,而是数据质量的保障。技术方案可以抄,但数据质量的坑只能一个个踩过来。建议刚开始做的时候,宁可功能少一点,也要把数据校验和监控做扎实,后面会省很多事。
最后分享一个小技巧:给每个数据表加一个数据质量看板,实时显示数据量、缺失率、异常值比例。这个看板能帮你提前发现90%的数据问题,比事后排查高效得多。