从行情API到稳定行情系统:量化交易数据接入的实践与坑
做量化和做交易工具的朋友有一个绕不开的起点拿到靠谱的行情数据。很多人一开始觉得调一套行情API很简单拿到token、请求一次、解析返回、打印出来齐活。但真到了要支撑一个实盘策略、或者要做一个多品种监控面板的时候各种问题就冒出来了行情断了没人知道、数据对不上、延迟不稳定、重启之后补数据补到头皮发麻。这篇文章不聊太高深的理论就说说从调通一个行情API到设计一套能用的行情系统这条路我实际踩过的一些坑和最后沉淀下来的做法希望对准备入坑或者正在重构行情模块的朋友有点参考价值。这个内容适合谁呢一类是刚接触量化交易、想自己拉行情做研究的新手另一类是已经在用第三方行情接口但觉得现在这套拉数据的方式不够稳、想重新设计一下的老手。不管你是做A股、期货还是数字货币核心的接口通信逻辑、数据存储设计、断线重连机制思路都是相通的。我会以常见的REST WebSocket组合为例来展开这部分无论是哪家行情服务商套路都差不多。1. 行情API到底解决什么问题1.1 行情数据的三种层级快照、Tick和K线行情数据听起来简单无非就是价格、成交量、时间。但真正做进系统里你会发现不同场景对数据的要求完全不同。我把行情数据大概分成三个层级这个划分决定了你的系统复杂度L1快照Level 1 Snapshot这是最常见的最新一档数据通常包含最新价、涨跌幅、成交量、成交额、买一卖一价和量。盯盘界面、简单的告警系统、分钟级策略用这个层级就够了。Tick逐笔成交/逐笔委托每一笔真实成交的价格和数量都记录下来了还有逐笔委托明细。这对高频策略、盘口分析和交易成本建模是必须的。数据量非常大一天的Tick数据往往是快照数据的几十上百倍。K线Bar按1秒、1分钟、5分钟、日线等周期聚合出来的数据。这个是最容易处理的因为它是经过聚合的数据量小、规律性强适合做回测和趋势判断。我见过不少人一上来就直接请求K线以为这就是全部了。可一旦你要做的是实盘监控或者高频回测没有Tick数据是很被动的。所以做系统设计之前先问自己一句我的策略或者业务到底需要哪一层的数据这个问题的答案直接决定了数据量级、存储成本和接口对接方式。1.2 自建采集与现成API的取舍有些团队喜欢自己写爬虫去交易所或者财经网站抓数据或者直接用开源库里的数据源。我的建议是做研究、做回测没问题做实盘系统还是老老实实接正规的行情API。原因很直接稳定性自建爬虫依赖网页结构别人前端改个class名你的解析就崩。而正规行情API是面向程序调用的有明确的版本控制和契约。数据准确性网页展示的数据经常是延迟或脱水的API的Tick数据则更接近原始成交。合规性很多行情服务商对数据转发、商用有明确授权要求自己爬的数据在法律和合同层面风险不可控。当然自建采集也不是一无是处比如你想采集一些冷门品种、小交易所的数据没有现成API可用的时候只能自己动手。我的经验是自建采集做补充来源可以但核心行情链路一定要有至少一个稳定、可追溯、有SLA保障的API通道不然生产环境一旦出问题就是连锁反应。2. 调通接口前先想清楚这几个关键参数2.1 接口类型选择REST与WebSocket怎么分工行情API大体的接口形式就两种理解清楚它们的定位后面的系统设计会很顺。REST接口一问一答请求一次拿到一个结果。适合做历史K线下载、启动时的初始化快照拉取或者低频的状态查询。它的特点是简单、天然支持重试但拿不到连续更新的数据。WebSocket接口客户端与服务端保持一条长连接服务端有数据就主动推给客户端。适合实时行情流延迟低是实盘系统的核心通道。缺点是连接状态需要自己管理断线、心跳、补数据都要考虑。实际架构里这两种通常会配合使用启动时先通过REST拉一次全量快照或最近一段K线把底仓建好然后通过WebSocket订阅实时增量这样既能快速初始化又不会在行情量大的时候被推送打蒙。这种REST补底 WebSocket增量的模式是我目前在用的标准方案。2.2 鉴权方式与权限粒度各家API的鉴权方式大同小异一般是给一组API Key和Secret Key请求签名后带上。需要注意的不是签名本身而是权限粒度的问题。具体来说只读权限和交易权限要分开如果你的Key允许下单那请务必限制IP白名单。我见过有人把带交易权限的Key放进代码仓库的配置里泄露之后被恶意下单资金损失是真实发生过的。行情模块只需要只读权限。订阅权限与数据范围有的服务商按产品线分权限比如A股行情和期货行情是分开授权的。接入前先确认你的Key能订阅哪些市场免得调通了数据源上线时发现拿不到某些品种的数据。连接数限制很多服务商对同一个Key的并发连接数是有限制的。比如同一个Key只能建立5条WebSocket连接这直接决定了你系统里行情网关需要部署几个实例、每个实例负责哪些品种。2.3 关键参数时区、复权、时间戳精度这几个参数看起来小但出错的时候非常隐蔽而且是那种数据好像没问题但就是不对劲的类型我一个个说。时区行情API返回的时间服务器时间用的哪个时区一定要确认清楚。常见的坑是服务端用UTC你本地用的是UTC8然后你直接用字符串做比较结果差8小时。我的处理方法是全链路统一用Unix毫秒时间戳只在展示层转成本地时间。复权因子做历史K线回测前复权、后复权和不复权的数据差别巨大。除权除息日前后不复权的价格看起来像跳空了一样。这个不是API能替你决定的需要在系统里统一配置按哪种复权方式进行存储和计算。时间戳精度有的API返回秒级时间戳有的返回毫秒还有的是带时区的ISO字符串。我在系统里用一个统一的数据结构去接收然后在入口处强制转成毫秒整数避免下游所有模块各自做转换时出现不一致。3. 从一次成功请求到稳定行情流核心环节实现3.1 建立行情连接订阅与心跳这一步是整个系统的地基。我用一段Python伪代码说明一下最核心的连接逻辑这套模式基本适用于所有WebSocket行情接口。import asyncio import websockets import json import time import hmac import hashlib class QuoteClient: def __init__(self, api_key, api_secret, url): self.api_key api_key self.api_secret api_secret self.url url self.ws None self.subscribed set() self.last_pong time.time() async def connect(self): self.ws await websockets.connect(self.url, ping_intervalNone) # 不少服务商要求连接后先发一条鉴权消息 await self.auth() # 恢复订阅把之前订阅过的频道重新订阅一遍 for channel in self.subscribed: await self.subscribe(channel) async def auth(self): # 以某个服务商为例发送带签名的鉴权请求 timestamp str(int(time.time() * 1000)) sign_str f{timestamp}GET/ws.encode() sign hmac.new(self.api_secret.encode(), sign_str, hashlib.sha256).hexdigest() auth_msg { op: auth, args: [self.api_key, timestamp, sign] } await self.ws.send(json.dumps(auth_msg)) async def subscribe(self, channel): self.subscribed.add(channel) sub_msg {op: subscribe, args: [channel]} await self.ws.send(json.dumps(sub_msg)) async def heartbeat(self): while True: await asyncio.sleep(15) await self.ws.send(json.dumps({op: ping})) # 如果超过一定时间没收到pong主动断开重连 if time.time() - self.last_pong 30: await self.ws.close() break async def run(self): while True: try: await self.connect() # 并发执行消息接收和心跳 await asyncio.gather(self.receive(), self.heartbeat()) except Exception as e: print(f[行情连接异常] {e}) await asyncio.sleep(3) continue这里的核心细节有两个关闭自带的ping_interval很多WebSocket库默认会发RFC标准的ping帧。但行情服务商通常不认这个它们有自己的应用层心跳协议比如自定义的{op:ping}消息。如果你不做应用层心跳只靠TCP层面的心跳连接看似没断实际上服务端早就把连接回收了数据停了半天你都不知道。订阅集合与重连恢复连接断开后重连不只是把socket建起来还要把之前订阅的频道重新订阅一遍。我用一个set()把已订阅频道记录下来重连后逐个恢复这个操作必须做不然就会有连接正常但没数据的假象。3.2 数据解析与序列化行情消息推过来通常是一条JSON但如果你直接把它丢给下游让每个模块自己解析后面就乱了。我的原则是在入口统一解码转成内部结构体再对外提供类型化接口。举例来说一条Tick消息可能有这样的字段class Tick: symbol: str price: float volume: float amount: float direction: str # buy / sell trade_time_ms: int seq: int统一解析的意义在于让时间戳、价格精度、字段缺省值在下游完全一致。不要在A模块里用float(price)在B模块里又用Decimal(price)最后对账对不上。因为行情数据是全系统公用的基础数据它的解析规则直接影响后面所有策略模块和存储模块的一致性。存储这块如果数据量小直接写关系型数据库问题不大但如果做实盘Tick级存储数据量就上来了。我建议在内存里做一次批处理比如攒够100条或者每200毫秒刷一次盘写入库而不是来一条写一条。这样能极大减轻数据库压力也减少磁盘随机写带来的性能抖动。文件格式我推荐Parquet或者Arrow按日期分区存储回放和统计都很方便。3.3 断线重连与数据补拉断线重连只是第一步真正麻烦的是补数据。因为断线期间你可能漏掉了一段Tick直接恢复订阅之后中间那段数据是空缺的。这不是行情服务商的问题是所有基于推送机制的系统都要面对的推过来的数据不保证你能一直在线接收到。我现在的做法是维护了一个状态机大致流程是这样的检测到WebSocket断开记录断开时间t1。进入重连中状态每3秒尝试重连最多重试5次。重连成功后通过REST接口拉取t1到现在这段时间的Tick数据我这里用的是按时间增量拉取的方式避免重复拿之前已经入库的数据。数据补拉完成后再恢复实时订阅。如果重试5次都失败就发告警通知人工处理。这里有一个值得注意的点重连成功之后先补数据再订阅顺序不要反。如果你先订阅了实时增量又去补历史数据到达的顺序就会交错需要额外做去重合并很容易乱。4. 行情系统的整体架构设计4.1 单机脚本到分层架构最开始做行情接入很多人就写一个脚本跑起来之后把数据打成日志文件就完事了。如果只是自用研究这个方案完全可以。但一旦要接入实盘、要支撑多个策略脚本这种单机形态就不够用了。故障恢复、并发订阅、数据多路分发、权限隔离这些需求逼着你要做分层设计。我现在常用的分层架构是四层层次职责关键点接入层维护行情连接、鉴权、心跳只做连接和协议处理不做业务逻辑网关层解析消息、统一数据模型、消息路由是系统的数据规范入口核心服务层行情存储、数据补拉、策略订阅管理处理业务状态应用层行情页面、策略引擎、告警服务面向用户和交易每一层保持独立部署和独立扩展能力。比如接入层连接数不够了多部署一个实例把品种分组订阅存储压力大了核心服务层增加消息缓冲。这种横向扩展的能力是单机脚本完全给不了的。4.2 延迟指标怎么量化行情系统做得再好也要有数字说话。延迟这个指标不能只看服务商宣传的毫秒级推送要自己动手测。我测延迟用的是本地时钟对比法在行情消息到达的第一时间记录本地时间t_recv。同时把行情消息里自带的服务端时间戳t_server透传。延迟近似值 t_recv - t_server前提是本地时钟和服务端时钟做了NTP同步。测出来的延迟要做分位数统计不要只看平均值。平均值是会被极端值拉低的真正影响交易的是P90、P99这些尾部延迟。我见过一个系统平均延迟5毫秒但P99延迟300毫秒策略里的抢单逻辑偶尔就会踩到这个坑。所以上线前一定要把P50、P90、P99都统计出来。4.3 消息队列和存储的选型思路行情系统内部模块间通信不建议用HTTP反复调用吞吐不够延迟也高。一般会引入消息中间件做解耦。选型的时候有两条路线轻量级方案如果系统规模不大单机部署直接用进程内的消息队列就够了。比如Python里用asyncio.Queue每条消息一推一收简单直接。跨进程乃至跨机器方案如果网关和策略不在同一台机器上就需要引入真正的消息中间件。我个人在行情场景里更偏向Kafka。它基于拉取模型消费者能够自主控制处理节奏而基于推送模型的中间件在瞬时高并发下容易把消费端打崩。Kafka的日志保留机制也能顺便解决部分历史数据回放的问题。存储层要分开说近期数据比如最近7天的Tick存时序数据库或者Parquet文件重点在查询速度和回放效率。长期数据月度、年度的K线存列存格式做统计研究的时候批量扫描效率高。4.4 监控和告警行情系统的生命线一套行情系统最怕的不是延迟高而是你以为在跑其实数据早断了。所以监控是最后一道防线也是很多人容易漏掉的一环。我在生产环境里至少会盯这些指标最后一条消息的时间戳每个订阅通道超过30秒没收到新消息立刻告警。连接状态WebSocket连接断开、重连次数超过阈值需要通知值班人员。消息积压数量消费者处理速度跟不上消息生产速度时积压会越来越大这一步能提前发现。数据完整性通过比对REST接口拉取的当前快照和WebSocket收到的最近Tick检测有没有漏数据。告警方式我用的是企业微信机器人加一个简单的Webhook把告警信息推送到群里。重要的告警比如行情全断配置电话或短信。不夸张地说这套监控系统曾经在半夜帮我避免了一次无法挽回的损失值得花时间做。另外还有一个很容易被忽略的小细节行情系统的日志必须带消息序号或时间戳。排查问题的时候没有时序日志几乎等于没有日志因为你根本不知道异常是发生在哪一秒、和哪条订阅频道相关。5. 常见问题与排查技巧实录这里整理一下我做行情系统过程中遇到比较多的问题做成一个速查表方便遇到同类问题时对照排查。问题现象可能原因排查思路解决方案连接正常但收不到数据订阅没有成功或订阅频道名不对检查订阅后返回的success回执不只看请求发送成功对订阅响应做持久化记录确保收到success才算订阅成功行情数据偶尔跳变本地时间与服务端时间不同步对比NTP时间检查时间戳精度全链路统一用毫秒时间戳增加时间同步重启后行情数据对不上重启期间漏数据没有补拉机制检查写库日志看中间是否有断档实现断线重连后的增量补拉内存占用持续上升消息处理不及时队列被撑大观察消费速率和生产速率比值优化消费逻辑必要时增加消费者数量或引入消息队列解耦数据重复重连后重复订阅同一频道检查订阅恢复逻辑里set()去重是否生效统一维护订阅集合幂等订阅行情收盘后还有残留连接没有处理休市或非交易时段状态查看错误日志中是否有非交易时段的服务器异常推送在系统里加入交易时段判断非交易时段主动断开订阅5.1 订阅成功了但还是没数据怎么定位这是我第一次做实时行情就踩过的坑所以想单独多说一句。当时我以为订阅消息发出去了就万事大吉但实际上很多服务商的订阅协议是异步回执的就是你发出{op:subscribe,args:[trade.BTCUSDT]}之后服务端可能过一会儿才返回一条{op:subscribe,success:true}的消息。如果你只是发出去了、不检查回执那么遇到频道名写错、权限不足的情况你连数据都收不到完全找不到原因。后来的做法是维护一张订阅状态表key是频道名value是订阅状态。收到success回执才把状态标记为已订阅并开始计时接收数据。如果超过一定秒数没收到数据主动查询一次REST快照来确认这个频道是否有数据在产生。这个排查链路能帮你快速定位是频道本身没行情还是你的订阅没生效。5.2 行情数据入库的性能问题有一阵子我们系统接收的Tick量涨得很快数据库写入开始出现明显延迟。查下来发现每条Tick都走一次INSERT磁盘写放大严重。后来改成攒批写入情况立刻缓解。具体做法是用一个列表在内存里攒消息攒到500条或者每300毫秒定时刷一次先写一次批量写入。另外数据库字段设计上也要考虑行情数据的特点。比如symbol字段要建索引trade_time_ms要建索引查询才快。这两个字段是行情数据最常用来检索的条件。联合索引通常比单列索引更高效我建议把symbol trade_time_ms做成一个联合索引。用不用分布式数据库取决于你的数据量但表结构设计一定要提前考虑查询模式不然上线之后再做迁移非常痛苦。5.3 时区问题的经典坑最后说一下时区。这个问题我之所以单独拿出来讲是因为它不出显眼的问题但会在对账单、统计报表的时间轴上让你怀疑人生。常见场景是这样的数据库里存的是北京时间但接口返回的是UTC时间你直接拿这个字符串去做字符串比较或者直接当成北京时间传给前端展示于是K线图上最后两根K线突然消失或者统计结果每天差8小时。排查了很久才发现是时区问题。我的建议是从行情API入口就把时间转成统一的Unix毫秒时间戳然后存储、传输、计算全链路都用它。只有在最终展示的时候才根据用户的时区转成字符串。任何模块都不要自己再去解析时间字符串时间处理逻辑只保留在唯一的入口和出口。6. 从实践沉淀下来的一些方法论项目做完之后回看整个从接口调通到系统设计的过程我觉得有几点经验值得分享给正在做类似工作的朋友。第一先用最小闭环跑通再做架构。很多人容易犯的毛病是系统设计先画了一堆模块图结果连API的一次请求都没调通过就开始写框架代码。我的顺序是先把一个品种、一类数据、一条链路的完整流程跑通比如从交易所API拿到Tick解析完写进数据库再在监控界面上显示出来。这个闭环一旦通了系统的各个模块边界也就自然清楚了这时候再做架构的抽象和扩展质量会高很多。第二行情系统里可观测性要比性能优化更优先。数据断了、延迟高了、重复了这些问题的确都要解决但如果连有没有在正常跑都看不清楚一切优化都是盲人摸象。所以我建议在系统第一天就加上监控和日志不要等到出问题了再补。第三对外部API永远保持不信任。即使服务商承诺了SLA你也要假设网络会断、服务会挂、数据会缺。所以要有一套独立的校验手段比如定期用REST接口拉快照来核对WebSocket推送数据确保数据的连续性和正确性。这是生产系统稳定运行的底线思维。行情API本身只是一个通道真正决定系统质量的是连接管理、数据工程、异常处理这些细节。把这些细节做扎实之后你会发现后面接新数据源、加新策略都顺很多。做行情系统是个体力活但也是一件越做越有底气的事。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →