量化交易实时行情接入:从API到策略层的工程实践 我最早做量化交易的时候天真地以为接行情就是把 WebSocket 连上回调里把 JSON 拿出来json.loads之后往策略里一扔就完事。等真正上了五档行情Level-2之后才发现完全不是那么回事每秒几十笔的 tick 数据还好但五档快照叠加逐笔委托数据量直接上了一个数量级策略里稍不注意延迟就从毫秒级拖到百毫秒级更麻烦的是断线重连、数据乱序、扩容时订单簿重建成本任何一个没处理干净上线之后都是事故。这篇文章我就结合自己的实战经验把从行情 API 到策略层这段路的工程设计完整拆一遍。这篇文章适合已经在用 Python 写策略、但还没系统处理过实时深度行情的同学也适合准备把策略从分钟级降到 tick 级、正犹豫要不要接五档的人。我会从行情数据的真实形态讲起再说 API 选型和接入层该做的活然后把数据管线的设计思路、策略层消费行情的方式逐步展开最后分享我踩过的几个坑。代码用 Python 写整体思路对所有语言都通用。1. 先说清楚行情数据的真实形态Tick、快照与五档的边界1.1 行情不是只有一种Level-1、Level-2 和全推的区别很多新手把实时行情理解成一个统一的东西实际上量化场景里至少要区分三层基础快照、逐笔成交、逐笔委托。日常手机 App 里看到的最新价涨跌幅买一卖一属于 Level-1 快照通常是每秒或每 3 秒推一次数据量很小做分钟级策略完全够用。五档行情一般指 Level-2 快照里带上买卖各五档的挂单价格和数量更新频率更高通常 1 秒甚至 500 毫秒一次。再往上还有十档、四十档以及逐笔成交和逐笔委托——后者才是真正的全推送每一笔真实发生的委托和成交都会推给你数据量按每秒几百到几千条计。我做日内高频相关的策略时五档快照是底线逐笔委托能拿到的话做订单簿失衡类信号会准很多。但这里要有个清醒的认知数据量不是免费午餐。逐笔委托的全推送高峰期一秒几千条消息很常见Python 如果处理不好单个消息的解码和分发CPU 很容易被打满。所以先明确你需要哪个层级的数据再决定后面的架构而不是行情源给你什么你就接什么。1.2 五档行情在策略里的真实用途不只是看一眼买卖盘五档数据在策略层最常见的应用有三个。第一是订单簿失衡信号计算买一到买五的总量和卖一到卖五的总量的比值比值极端时往往代表短期买卖力量不平衡可以用来做超短线入场辅助。第二是滑点估计和下单执行优化当你准备下大单时光看最新价是不够的必须知道吃穿到第几档才能全部成交这直接决定你的限价单挂在哪一档。第三是撮合模拟回测和盘前仿真时用五档数据模拟你下单后能成交的价格和数量比按收盘价成交真实得多。这三个用途对应不同的工程要求。做订单簿失衡信号需要维护一个能快速求和的五档深度结构做执行优化需要精确到每个价格档位的累计数量做撮合模拟则需要考虑你的订单改变订单簿之后的后续状态。很多策略框架只处理 tick 价格流对五档数据只取了买卖一档等于把最值钱的信息浪费了大半。1.3 一个被普遍低估的工程问题快照与增量接了五档行情之后第一个遇到的实际问题往往是行情源同时给你全量快照和增量更新两种消息。全量快照是当前时刻完整的一到五档增量更新则是某一档价格变了、数量改了或者某档位撤单清空。正确做法是用全量快照做基准之后应用增量流不断更新。难点在于快照和增量的版本要对齐——如果增量比快照早到或者中间断线恢复后漏了若干条增量订单簿就会悄悄错乱。这个问题在工程上的解法后面第 4 节会详细展开。这里先记住一个结论五档行情的接入不是一个来什么存什么的问题而是一个典型的快照 增量 序列号对齐问题。不解决它你的订单簿就是一本可能翻错页的账本。2. 数据 API 选型从券商柜台到第三方推送接口边界怎么划2.1 国内主流行情源的大致谱系Python 做量化数据源我把它粗略分成三类。第一类是互联网公共接口或半公开接口比如常见的数据平台提供的 HTTP/WebSocket 接口优点是上手快、免费或低成本缺点是延迟偏高、历史数据与实时数据可能会有一致性偏差通常适合研究、回测和低频实盘。第二类是券商或期货公司提供的官方接口比如期货领域常用的 CTP 行情接口这类接口延迟低、数据质量有保障但通常需要开通相应账户权限API 风格也更底层Python 往往得通过封装层来调用。第三类是加密货币交易所的公开 WebSocket 行情这类接口标准化程度高、文档清晰、免费很适合用来练习实时行情处理和策略回测但市场微观结构和 A 股/期货差异很大不能直接套用。从工程设计角度看我不建议在策略代码里直接绑死某一家 API。理由很简单接口会变、延迟会变、甚至你用的行情商哪天停了也不奇怪。所以选型第一步不是比较谁家的字段更多而是确认它的鉴权方式、消息格式、断线语义、订阅模型这几个决定架构的要素。2.2 我用下来最看重的四个选型标准第一个是鉴权与订阅模型。有些接口用 Token 鉴权每连接一次需要重新申请有些接口用账号密码签名过期时间很长。这会直接影响重连逻辑的复杂度。订阅模型也要看是连接后一次性订阅全部还是可以按需增量订阅、取消订阅。支持增量订阅的接口在策略需要盘中切换品种时会舒服很多。第二个是消息格式与序列化。JSON 最常见调试方便但解析开销大二进制协议如 protobuf、FlatBuffers 或自定义二进制解析快、体积小但调试麻烦。我自己的原则是研究阶段用 JSON 没问题延迟敏感的实盘系统尽量让行情源直接给你二进制原始数据Python 里再用结构化解包能省掉大量 JSON parse 的 CPU。第三个是推送频率与数据完整性标记。确认接口推送频率、是否有每次推送的序号、是否有快照/增量类型标记、断线重连后如何重新同步这些比文档里漂亮的示例代码重要得多。第四个是限流和最大订阅数。很多公共接口虽然写着不限订阅但实际连接数或请求频率超过阈值就会被断开甚至封禁。接入前先压测一下你计划中的订阅数量是否在安全区间别等上线之后被限流打击。我用一个表格总结一下选型时的关注点方便你对照自己手上的 API维度需要确认的问题影响鉴权方式Token 有效期续期机制决定重连后是否需要重新登录数据格式JSON / 二进制 / protobuf决定解析性能与调试成本推送频率固定周期 or 变化频率是否有最大频率限制决定消费端是否要做削峰序列号每条消息是否有自增序号决定乱序检测和增量对齐方案快照与增量是否都有如何区分以什么作为触发重置的基准决定订单簿重建的正确性断线语义重连后是全量重推还是可以续传决定重连恢复时间2.3 无论选哪家先做一层接口抽象无论最后用哪家数据源我强烈建议写一个MarketDataClient抽象基类定义好connect()、subscribe(symbols)、on_message()、close()这些标准方法再把具体厂商的实现放到子类里。这样策略层只依赖抽象接口数据源切换就不需要改策略代码。看起来多写了几行但当你从公共接口换到券商 CTP 或者从某一家换到另一家时会无比庆幸当时做了这层隔离。这里还要提醒一句别把行情源自己的数据结构直接传到策略层。比如 CTP 里的DepthMarketData字段风格、公共接口里的bids/asks数组结构都不一样最好在自己的数据管线里统一转换成你定义的标准模型。也就是说接口抽象解决的是厂商切换统一模型解决的是数据结构漂移两层都要做。3. 行情接入层的工程细节解码、校验、断线重连、心跳保活3.1 一个最小可用的 WebSocket 行情客户端在很多场景下行情接入最常见的形式就是 WebSocket。我写过一个最小客户端骨架去掉业务逻辑之后大概长这样import asyncio import json import websockets from typing import Optional, Callable, Awaitable class StreamClient: def __init__(self, url: str, token: str, on_message: Callable[[dict], Awaitable[None]]): self.url url self.token token self.on_message on_message self.ws: Optional[websockets.WebSocketClientProtocol] None self._running False self._last_seq: Optional[int] None async def connect(self) - None: self.ws await websockets.connect(self.url) await self.ws.send(json.dumps({ op: auth, token: self.token })) self._running True # 订阅放到外部由具体策略决定 async for raw in self.ws: msg json.loads(raw) if msg.get(type) in (ping, pong): continue await self.on_message(msg) async def run_forever(self) - None: while self._running: try: await self.connect() except (websockets.ConnectionClosed, asyncio.TimeoutError, OSError) as exc: print(fconnection error: {exc}, retry in 3s) await asyncio.sleep(3) async def close(self) - None: self._running False if self.ws: await self.ws.close()这个骨架的核心思路是连接与业务解耦消息统一走on_message回调断线时整个连接循环都由run_forever负责重试。单看代码很简单但这里藏着一个经验重连时不能只重连而不重新订阅和重新鉴权。很多接口在连接建立后需要重新发订阅指令部分接口还需要在重连后主动发起一次快照请求。所以我的习惯是connect()内部按顺序完成建立连接 - 鉴权 - 订阅 - 开始消费消息任何一步失败都抛异常让上层重连。3.2 心跳与重连的正确姿势幂等订阅和序列号校验公共接口通常有服务端心跳比如每隔 15 秒发一个ping客户端只需要回pong即可。但有些自建的行情服务不会主动发心跳需要客户端定时发ping超过一定时间没收到响应就要主动断开重连。这个保活逻辑我建议独立成一个协程和时间驱动的重连机制配合使用。更关键的是重连后的状态恢复。第一次连接后你订阅了 50 个合约断线重连后如果订阅指令是重复发的一般没问题因为现网接口基本都支持幂等订阅但你必须记录我当前订阅了哪些合约重连成功后再逐批订阅不能依赖服务端记住旧连接的状态。这一点对自建网关尤其重要因为服务端重启后所有连接状态都丢了你以为还订阅着实际服务端已经把你忘了。3.3 消息解码与性能json.loads不是唯一解如果你做分钟级策略json.loads完全够用。但五档行情一上来尤其盘中高峰期每个 tick 都携带买卖十档的价格和数量JSON 解析的 CPU 开销会非常明显。我在压测里遇到过单路行情每秒 1000 条推送、单条消息大小 2KB 的场景用纯json.loads在普通云主机上解析就要吃掉一个核的一半。解法有几个方向如果 API 支持二进制协议直接用二进制如果只有 JSON可以先做一次 profile看热点是不是在解析上如果解析确实是瓶颈可以考虑用orjson这种高性能库替换标准库实测通常能提升 2 到 5 倍。另一个容易被忽略的点是避免在回调里做重活回调只负责解码和投递真正的数据处理放到下游。回调里哪怕只是打日志在高频下都会成为瓶颈。4. 从原始消息到策略事件中间那层数据管线怎么设计4.1 为什么不能直接把行情 JSON 丢给策略我见过不少项目策略函数签名直接是strategy(price: float, bids: list, asks: list)看着很直观但盘中一旦发现数据乱序、涨跌停、停牌复牌或者其他数据异常根本无从查起。因为 JSON 只是当时的原始报文不是一个可以直接做逻辑判断的一致性视图。所以在行情源和策略之间必须加一层数据管线负责把原始报文解码成结构化的领域事件同时做校验、清洗、序列号检查。这层管线就是从数据 API 到策略层的工程设计里最核心的部分。它做的事情可以概括为三件转译——把厂商报文转成内部统一的OrderBookUpdate、TradeTick、Ticker事件对齐——保证增量更新基于正确的快照分发——把事件推给关心它的策略模块。4.2 统一事件模型用 dataclass 定义你的领域对象用dataclass定义领域模型是我在 Python 项目里最喜欢的做法。它既轻量又能保证字段可读性还可以嵌套表达复合结构。下面是我常用的一组模型from dataclasses import dataclass, field from typing import List, Optional dataclass class PriceLevel: price: float size: float order_count: Optional[int] None dataclass class OrderBook: symbol: str ts: int # 交易所时间戳毫秒 bids: List[PriceLevel] # 买盘按价格从高到低 asks: List[PriceLevel] # 卖盘按价格从低到高 snap_seq: Optional[int] None # 快照序列号用于对齐 dataclass class TradeTick: symbol: str ts: int price: float size: float side: str # B / S / N dataclass class Ticker: symbol: str ts: int last: float open: float high: float low: float volume: float统一事件模型的价值在于策略层只需要认识这几个类不需要关心它们来自哪个行情源。换数据源时只需要在接入层写一个新的适配器把厂商报文转换成你的OrderBook和TradeTick策略一行都不用改。我记得第一次把策略从公共行情源迁移到另一家行情源时只改了一个适配器文件半小时就完成了切换。4.3 订单簿重建增量更新五档的核心逻辑五档增量的更新模式通常是给定一个档位设置新的价格和数量或者把某个档位清空。订单簿重建的目标就是维护一个始终与交易所状态一致的买卖深度表。下面是我常用的一种实现思路按价格水平组织成字典更新时直接覆写再定期重新排序生成列表from sortedcontainers import SortedDict class OrderBookManager: def __init__(self, levels: int 5): self.levels levels self._bids: SortedDict SortedDict() # price - size, 降序 self._asks: SortedDict SortedDict() # price - size, 升序 def apply_snapshot(self, book: OrderBook) - None: self._bids.clear() self._asks.clear() for level in book.bids: if level.size 0: self._bids[level.price] level.size for level in book.asks: if level.size 0: self._asks[level.price] level.size def apply_update(self, side: str, price: float, size: float) - None: book self._bids if side bid else self._asks if size 0: book.pop(price, None) else: book[price] size def snapshot(self) - OrderBook: bids [PriceLevel(p, self._bids[p]) for p in reversed(self._bids.keys()[-self.levels:])] # 最大的5档 asks [PriceLevel(p, self._asks[p]) for p in self._asks.keys()[:self.levels]] return OrderBook(symbol, ts0, bidsbids, asksasks)注意里面一个细节买盘价格需要从高到低排列卖盘价格从低到高排列这并不只是显示顺序问题而是策略里算买卖失衡时会依赖索引位置。SortedDict的keys()是升序所以买盘要reversed(...)取最后一小段再反转卖盘取前几个。这段逻辑很容易写错我在这上面吃过亏。真正的工程难点在于增量与快照的对齐。我的建议是给订单一簿维护一个last_seq字段行情源推送任何快照或增量时都带一个自增序号。如果接收到的增量序号不等于last_seq 1说明中间漏了数据这时候不能再继续应用增量而应该丢弃当前订单簿、请求一次全量快照重新建基准。整个过程要在毫秒级完成否则订单簿在盘中会长期处于不可信状态。4.4 线程模型与背压处理单线程事件循环 队列 快照兜底行情接入和策略运行通常分属不同类型的任务。我建议整体架构采用单线程事件循环行情源回调只负责把解码后的原始事件放入一个无界队列或内存管道策略层在同一个事件循环里从队列取事件、更新订单簿、运行信号、生成订单。这样天然避免了多线程共享订单簿的锁竞争问题也省掉了大量并发 bug。但单线程模型有一个风险如果队列消费者处理不过来内存会一直上涨。所以在生产环境里我一般会对队列长度做监控超过阈值就报警再设置一个兜底策略如果处理延迟超过 500ms就丢弃部分低频快照事件、只保留逐笔成交和最后一张快照确保核心信号不至于完全停滞。实盘系统里卡死往往比丢一点数据更可怕宁可主动丢弃低优先级消息也不要让整个系统崩溃。这条经验不经历一次盘中延迟飙升很难体会。5. 策略层如何消费行情信号计算、撮合模拟与风控联动5.1 以订单簿失衡为例五档数据怎么变成交易信号有了干净统一的OrderBook对象之后策略层的事情就变得单纯了。以订单簿失衡信号为例核心逻辑就几行def orderbook_imbalance(book: OrderBook) - float: bid_vol sum(level.size for level in book.bids) ask_vol sum(level.size for level in book.asks) total bid_vol ask_vol if total 0: return 0.0 return (bid_vol - ask_vol) / total数值接近 1 说明买盘压力大接近 -1 说明卖盘压力大。但实际策略里很少单看一个时刻的失衡值通常会和过去 N 个时刻的均值做对比或者结合价格变化方向过滤假信号。比如价格在上行但失衡值由正转负往往意味着推动力减弱这时候反而要小心。这些信号逻辑并不复杂真正决定盈亏的是对数据流的时序把握。五档行情的接入对信号的影响还有一个维度是延迟。同一个失衡信号如果拿到的是 200ms 前的快照和拿到的是最新快照交易结果可能完全不同。所以在策略层一定要给每个事件打上交易所时间戳和本地接收时间戳两个字段统计两者差值的分布如果延迟突然拉大很多信号都不能用。5.2 事件驱动框架回调函数 vs 消息总线怎么选小规模策略可以直接用回调函数order_book_manager.on_update(callback)行情来了就调用回调。这种方法最简单适合单策略、少品种。但如果你同时跑多个策略、多个品种回调的注册顺序、异常隔离、策略间的数据共享都会变得很难管理。我的经验是当策略数量超过两三个或者涉及多品种同时交易时就用轻量消息总线。消息总线的思路也不复杂每个策略模块订阅它关心的事件类型和品种行情管线发布OrderBookUpdate事件总线负责分发。Python 里pyee、blinker都是可以快速落地的实现。用总线的好处是解耦新策略加进来不需要改老策略坏处是出问题时链路变长、不好排查。所以总线里的事件一定要带唯一的 event_id 和 seq方便根据日志还原全程。5.3 策略中处理延迟与乱序的防御式写法实时行情里乱序是无法避免的无论你用的是公共接口还是官方接口。一种常见的乱序是旧时间戳的消息晚于新时间戳到达。处理方法是不要盲目用新到的数据覆盖旧数据先判断事件里的交易所时间戳是否早于你订单簿里记录的时间戳是的话直接丢弃或只更新局部字段。另一个问题是盘中出现瞬间的大单跳变或价格异常比如某档价格突然从 10 块变成 0.01 或者变成 1e9。别看这种错误少一旦出现你的信号计算就会爆表。防御式写法是在更新订单簿前做一层校验价格大于 0、价格落在当日合理波动范围内、数量非负不满足的直接丢弃并记日志。这些校验代码很短但能避免绝大多数数据源抽风导致策略乱下单的极端案例。5.4 从信号到下单模拟撮合与滑点设计策略层生成信号之后不能直接拿去下单必须先过一个模拟撮合模块。这个模块的作用是如果当前五档卖一有 100 手你想买 200 手那你的成交价就是卖一价格成交 100 手、卖二价格成交 100 手平均成交价取决于穿过的档位。这就是五档数据在策略执行侧的直接价值。模拟撮合还能帮你估算滑点。实盘里你的市价单可能吃掉几档模拟时同样的逻辑走一遍记录理论成交价和信号触发价的差值长期累积下来就是策略的隐性成本。我见过不少策略回测好看、实盘赚不到钱就是因为忽略了五档深度对成交量的限制。把模拟撮合纳入策略层之后下单逻辑才真正和数据管线形成了闭环。6. 回测和实盘的行情差异以及我踩过的几个坑6.1 回测用的历史数据 vs 实盘实时数据回测时用的历史 tick 数据和实盘实时收到的行情有本质差异。历史 tick 数据往往是事后整理的可能做了清洗和去重盘中收到的实时数据可能包含错误价、重复消息、乱序消息。如果回测代码不做任何防御式校验跑得通不代表实盘也能跑反过来实盘里合理的丢弃逻辑必须提前在回测里测试不然你根本不知道丢弃了多少数据对结果有多大影响。所以我的习惯是同一个策略先用历史 tick 回测一遍再把相同的行情处理管线接到实时数据源上做盘前仿真观察订单簿是否稳定、信号是否异变、延迟是否有异常。仿真跑上几天确认各项指标和回测接近才考虑小仓位实盘。这个流程听着繁琐但能过滤掉大量数据源层面的坑。6.2 我踩过的几个具体的坑坑一停牌复牌导致五档数据突变。股票停牌期间五档可能没有正常更新复牌瞬间快照数据从异常状态跳变回正常状态。如果不做校验策略可能把复牌瞬间的价格跳跃当成正常机会下出意想不到的订单。我的做法是监控每笔 tick 的时间间隔和价格变化幅度超过阈值就强制丢弃并重新订阅全量快照。坑二API 限频被悄悄触发。有些公共 API 没有明确报错只是在你连续高频订阅/退订时慢慢把连接断掉。排查这种问题很难因为表面看一切正常只是行情断断续续。解决办法是给服务端请求频率加一个本地节流阀每次订阅/退订都记录时间戳单位时间内的操作数超过安全值就排队等待。坑三时间戳时区处理不一致。同一份数据里有的接口用 Asia/Shanghai 时区的本地字符串有的用 UTC 时间戳甚至有的用 UNIX 毫秒。如果策略直接拿这些时间戳做计算回测和实盘对不上就是家常便饭。现在我的规矩是所有时间在进入数据管线的一刻就统一转成上海时区的毫秒时间戳存整数后续任何运算都基于整数毫秒避免时区混淆。坑四行情暂停时该不该输出信号。上午休市、下午开盘、临时停牌这些状态下五档数据可能长时间不更新。如果策略在数据静止时仍然使用旧数据计算信号就可能产生无效交易。我的处理是维护一个data_freshness字段距离上一次有效更新时间超过 N 秒就自动把信号置为无效。实盘里这个字段救过我很多次。6.3 工程化兜底手段健康度监控、降级与日志数据管线的整体健康度必须可视化。我每台行情机上都跑一个最小监控面板展示四个数字当前队列长度、最近一条行情延迟交易所时间戳与接收时间戳之差、最近一分钟消息条数、订单簿重建重置次数。订单簿重置次数尤其关键如果短时间内重置超过几次说明增量数据连续性有问题必须及时排查。日志方面行情数据量太大不建议全量记日志。我一般只记录异常事件和降级事件比如乱序丢弃、序列号跳变、快照重建、队列溢出。正常行情数据的样本另存按分钟压缩归档方便事后复盘。最后强调一点所有日志里都要带上symbol、ts、seq三个字段否则出了问题根本没法按时间线回溯。这条规则看似基础我见过不止一个项目因为日志缺字段导致排查一个线上问题花了好几天。写在最后我自己的体会是实时行情接入这套东西本质上不是调 API的问题而是把不可靠的实时数据流转换成可靠策略事件的过程。从选型时给接口边界做抽象到接入层处理心跳重连和序列号对齐再到数据管线里的订单簿重建和背压控制每一层都在解决一个具体的数据工程问题。Python 在性能上不是最强的但用它搭一套行情接入和数据管线的原型配合良好的工程结构完全能满足中小规模的量化实战需求。最后再分享一个小技巧行情接入层写完以后先做一次 7x24 小时的模拟盘运行看看断线重连和重启恢复是否稳定再考虑接真钱。这一步省下来的可能是一整年的后悔。