ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

实时行情与五档行情接入量化策略:从API选型到事件驱动架构

实时行情与五档行情接入量化策略:从API选型到事件驱动架构 做量化这几年我有个特别深的体会策略代码写得好的人不少真正能稳定跑在资金上的却不多。绝大多数问题不出在策略本身的数学逻辑上而出在数据和工程之间那条看不见的缝里。尤其是实时行情和五档行情很多新手直接写一个循环去请求接口或者把收到的每一笔推送都在回调线程里算策略结果不是漏数据就是卡死轻则信号慢半拍重则整个程序直接崩掉。这篇文章就围绕一个核心问题展开实时行情和五档行情到底怎么接进策略层我会从数据 API 的选型聊到工程层的数据流设计再拆到策略引擎如何消费行情把“数据怎么来”“数据怎么存”“数据怎么用”这条链路完整讲清楚。内容会比较细涉及 WebSocket 连接管理、队列设计、行情快照与增量处理、事件驱动架构等前后端概念都有。如果你正在做自己的量化交易系统或者正在为团队搭建行情接入层这篇文章应该能帮你少走很多弯路。1. 行情数据体系搞懂实时行情和五档行情在策略中到底怎么用1.1 行情的三层结构Tick、快照、逐笔到底看什么很多刚接触量化的人会把行情理解成“K线”但实盘策略需要的数据颗粒度比K线细得多。通常行情体系里会区分三个层次。第一层是快照行情也叫盘口快照反映某个时刻的委托簿状态包括买卖五档或十档的价格和数量、最新成交价、今开、昨收、最高、最低、成交量、成交额等。这是策略最常用的行情形态。大多数 API 如无特殊说明返回的“实时行情”就是这种。快照通常是高频推送行情波动剧烈时每秒可能推好几笔平静时可能一两秒才推一笔。第二层是逐笔成交也叫 Tick-by-Tick是一笔一笔的真实成交记录包含成交时间、价格、数量、成交方向主动买还是主动卖等。逐笔成交比快照更细可以精确还原资金流动态。第三层是逐笔委托包括每一笔挂单、撤单、成交导致的委托簿变化。这一层是 Level-2 行情中数据量最大的部分通常需要单独付费购买对策略的提升也最直接因为它能揭示主力的挂单意图和撤单行为。标题里提到的“五档行情”在传统语境下指买卖各五档的盘口快照很多免费行情源提供的就是这个。如果是 Level-2 深度行情可能扩展到十档甚至全档位再叠加逐笔成交和逐笔委托。对个人量化玩家来说五档快照加逐笔成交已经能做出很多有效策略做市商或高频场景才真正需要完整逐笔委托。1.2 五档行情在策略中的真实价值不是只看“三买三卖”我见过不少人把五档行情拿来画个图就完了其实五档数据的核心价值在于它可以计算盘口失衡、委托压力、大单异动等结构特征。比如你能从买一到买五的总量对比卖一到卖五的总量算出委托失衡比你可以观察买一、买二的价格挂单是否在持续加厚或偷偷撤单你还能用逐笔成交的主动买卖方向识别大单是在吸筹还是出货。一个简单的例子买一档挂单突然从 500 手增加到 5000 手价格也顺势上抬而卖一档数量没有明显变化。这种盘口变化对短线策略来说是强特征但只看每分钟一次的K线根本感知不到。再比如五档价差结构当卖二到卖五档位连续稀疏只有卖一档厚厚盯着说明上方压力极有可能在某价位集中释放。这些信号全部来自五档快照区别在于你能否把数据转换成策略能用的特征。我在设计自己的系统时会把五档快照里的信息分拆成三类特征价格类特征、委托量类特征、委托变化速率类特征。价格类包括一档买卖价差、加权买卖价、盘口价格重心委托量类包括买卖总量比、单档密集度、撤单率估计变化速率类包括买一增量频率、卖三到卖五总量变化斜率。这些特征会在策略引擎里实时计算然后作为信号输入。1.3 行情的两种数据形态推送与拉取搞错会出大事行情数据的接入方式分两类推送和拉取选择直接影响你的系统设计。拉取方式最典型的就是 REST API你发起请求服务端返回当前快照。拉取的优点是开发简单调试直观但缺点很致命——你永远只能拿到“上一次请求瞬间”的数据两次请求之间发生什么你完全不知道。而且高频轮询会触发数据服务商的限流还可能被封 IP。所以 REST 适合做低频的估值、盘前准备、历史数据回放不适合做实时策略。推送方式的核心是 WebSocket服务端主动把行情推给你你被动接收。这种方式实时性最高只要服务端有行情就会立刻送到省去了轮询的心跳开销也不存在“照顾不到所有股票”的问题。所以目前的实时行情接入基本都走 WebSocket 通道。需要留意的是WebSocket 连接本身有断线风险服务器负载高时可能宕机网络环境也会影响稳定性所以工程上必须处理断线重连、数据补偿、心跳保活等一系列问题。这里多说一句推送行情并不代表“自动保证数据完整”。因为网络闪断、服务端重启、连接超时都会导致数据流中断你如果没做校验和补偿策略层拿到的数据就是残缺的后面算出来的指标全部不可信。这个在第五部分我会详细讲。2. 数据 API 接入从底层接口到工程封装2.1 行情源选型别盲目追“免费”先搞清楚业务需求行情数据源的选择是整套系统里最容易被低估的一环。不同数据源在实时性、稳定性、数据完整度、API 体验上差异极大而且各有各的限制。如果你只跑日线级策略免费数据源或者交易软件自带的导出行情就够了但如果你做日内策略尤其是秒级甚至 tick 级策略行情源的质量直接决定了策略能跑多深。选行情源时我会按下面几个维度打分第一是数据的原始颗粒度是只有快照还是有逐笔成交历史数据能不能回补第二是 API 的并发能力比如同时订阅多少只标的不会出现限流或延迟第三是稳定性有没有 SLA 保证是否经常停机维护第四是鉴权机制和数据版权限制避免商用后惹出麻烦第五是成本按量计费还是包月包年。针对个人和小团队我建议先确认自己做多高频。如果是秒级低频很多行情源都够用如果做的是真正的毫秒级盘口策略那基本只有交易所直连或券商极速柜台加 Level-2 行情源能满足。这两者的成本、开发难度、机房要求完全不是一个量级。别一开始就奔着最高配置去容易被工程复杂度拖死。2.2 接入实时行情 API 的完整代码骨架连接、订阅、回调、重连不管用哪家行情源WebSocket 接入的核心骨架都是类似的。我用 Python 写过一个封装示例逻辑是初始化连接、设置心跳、注册订阅、接收消息后分发到不同处理器。下面是一段简化过的代码展示了主循环和事件分发结构。import asyncio import json import websockets class MarketDataClient: def __init__(self, url, token, handlersNone): self.url url self.token token self.handlers handlers or {} # {snapshot: handle_snapshot, trade: handle_trade} self.ws None async def connect(self): # 连接并鉴权然后进入消息监听循环 async for websocket in websockets.connect(self.url): try: self.ws websocket await self._auth() await self._listen() except websockets.ConnectionClosed: # 断线后自动重连带指数退避 await asyncio.sleep(2) continue async def _auth(self): msg {cmd: auth, token: self.token} await self.ws.send(json.dumps(msg)) async def _listen(self): async for raw in self.ws: data json.loads(raw) msg_type data.get(type, snapshot) handler self.handlers.get(msg_type) if handler: # 注意这里不能 await 长任务要丢进队列或异步任务池 await handler(data)这段骨架里最需要注意的地方在_listen这一段。如果handler里做了耗时操作比如计算指标、写数据库、甚至触发下单整个 WebSocket 的接收循环就会被卡住。行情的推送频率非常高一个耗时操作没结束下一批行情已经在内存里排队连续几次卡顿就会出现严重的数据延迟。所以生产环境里我都是把接收和数据处理分开接收回调只做最轻量的解析然后立刻把数据放到队列里交给独立的工作线程处理。2.3 数据模型设计写字段名之前先想清楚时间戳的坑行情数据模型设计得好不好直接决定后续策略代码的复杂度。一个常见的错误是把 API 返回的字段名直接当一个字典到处传结果策略代码里到处是data[lastPrice]、data[bidVolumes][0]维护起来非常痛苦。我建议在每个行情源之上加一层数据模型层把原始消息解析成统一结构。比如定义一个MarketSnapshot数据类包含时间戳、代码、最新价、成交量、买卖各五档价格和数量等字段。这样策略层面对的是稳定对象行情源换掉后只需改模型转换层不用动策略逻辑。from dataclasses import dataclass from typing import List dataclass class MarketSnapshot: symbol: str exchange: str timestamp: int # 交易所时间毫秒 local_receive_time: int # 本地接收时间毫秒 last_price: float last_volume: int bid_prices: List[float] # 买一~买五 bid_volumes: List[int] ask_prices: List[float] # 卖一~卖五 ask_volumes: List[int]这个模型里我特意放了两个时间戳一个是timestamp代表交易所撮合时的时间另一个是local_receive_time代表我们自己系统收到这笔数据的本地时间。为什么要两个因为在分布式和网络环境下行情从交易所到你的程序之间是有延迟的交易所时间反映的是业务发生时刻本地接收时间反映的是业务到达时刻。分析系统延迟时必须用这两个时间戳做差值才定位得准是网络慢、源端慢、还是我们处理得慢。3. 实时数据管道别让行情把策略线程压垮3.1 为什么不能在行情回调里直接算策略很多第一次搭系统的同学都会犯同一个错误在 WebSocket 回调函数里顺手就把指标算了条件满足了下单日志也写了看起来一气呵成。但行情源是毫秒级推送如果你的指标计算是 O(n²) 复杂度或者里面有个网络请求回调耗时就会从 1 毫秒膨胀到几百毫秒。这时候新的行情不断进来回调堆积内存暴涨最后程序直接卡死或被杀。即使你的指标计算很快也可能出现另一个问题下单动作里包含券商接口调用或者数据库写入这些 IO 操作是不可控的遇到券商系统慢或网络抖动一次下单可能耗时几秒。在回调里做这种事等于把整个行情通道的风险都系在了一条链路上。正确的做法只有一种回调只做数据清洗和投递把真正的业务逻辑放进消费端。我把这种设计叫做“生产者和消费者分离”。生产者是行情连接负责接收原始数据消费者是策略引擎负责吃数据、算指标、发信号。中间解耦的介质是队列市场行情的“生产者”永远不知道“消费者”到底花了多少时间。3.2 用 Queue 还是用消息中间件看你的场景Python 里最简单的队列是queue.Queue它适合单机单进程的实时行情消费。生产者在回调里put消费者的工作线程get代码简单直观。如果你只是跑单一市场、单一策略queue.Queue完全够用。import queue import threading # 行情队列最大长度 5000避免无限堆积导致内存溢出 tick_queue queue.Queue(maxsize5000) def ws_handler(data): try: tick_queue.put_nowait(parse_snapshot(data)) except queue.Full: # 队列满了说明消费速度跟不上必须告警或降级 log_error(tick_queue full, dropping data) def strategy_worker(): while True: snapshot tick_queue.get() run_strategy(snapshot)注意put_nowait的用法。队列设了上限满了就不再阻塞直接记录日志后丢弃。这种“丢弃”策略看起来不好但比“无脑阻塞”要安全得多。因为队列一旦阻塞生产者那边行情就会越积越多最终把整个连接卡死。真实交易中行情短时间爆发是常态宁可主动丢弃超压的数据也不能让系统崩溃。当然丢弃之后要立刻告警说明消费能力有瓶颈。如果你的系统要支持多个策略同时订阅不同行情、或者需要跨进程共享行情数据那queue.Queue就不够了需要引入 Redis、Kafka 或者专业的消息中间件。这里的基本思路是行情源接收后先入中心化队列策略进程从队列里订阅自己关心的标的。加了一层中间件之后系统扩展性大幅提高但也引入了额外的延迟和运维复杂度个人项目要谨慎评估。3.3 乱序、重复、丢失实时行情工程的三大敌人实时行情接入最让人头疼的就是数据会出现“乱序、重复、丢失”这三种异常。先说乱序。同一只股票的行情在网络上可能走不同的路径到达你本地的顺序不一定和发生顺序一致。比如逐笔成交里先成交的一笔可能在网络中走得慢后到你的程序里。如果策略直接按接收顺序计算顺序就错了。解决办法是给每条行情带上序号或按时间戳做缓冲排序但这需要行情源支持。再说重复。网络重连后数据源为了补偿你断线期间的内容可能重复发送之前的行情。如果你不做去重成交量就会双倍计算策略看到的成交量和实际严重偏离。去重需要为每条行情做唯一键通常是“行情类型 代码 交易所时间 序号”。收到数据后先查这个键是否已处理是则跳过。最后是丢失。最常见的原因是断线重连的窗口期也可能来自我们自己的队列溢出。丢失很难完全避免关键是能否及时发现。做法是在每条行情里带一个自增序号连续两条之间的序号如果有跳跃说明中间漏了数据。发现漏数据后可以根据行情源的能力选择回补或重新订阅。如果是低频策略漏一两笔可能问题不大如果是高频策略漏任何一笔都可能导致信号失真这时就得考虑切换更稳定的数据通道。4. 策略层接入事件驱动架构是怎么落地的4.1 策略引擎的数据流从行情事件到策略上下文前面搭好了行情接入和数据管道接下来就是最关键的一环策略层怎么写。量化策略的实盘架构主流是事件驱动模型系统不在一个固定的“主循环”里轮询数据而是被动接收外部事件每个事件触发相应的策略代码。针对行情驱动型策略我比较推荐“上下文 事件回调”的结构。每个策略持有自己的状态上下文实时行情来了以后策略引擎先更新上下文然后让所有订阅了该标的的策略逻辑逐一判断是否要发出信号。这种模式的好处是策略之间天然隔离互不干扰缺点是事件分发需要精心设计不能让一个策略挂掉影响其他策略。下面是一个简化的策略引擎骨架展示了如何把行情数据分发到多个策略实例。class StrategyEngine: def __init__(self): self.strategies {} def register(self, symbol, strategy): self.strategies.setdefault(symbol, []).append(strategy) def on_snapshot(self, snapshot): # 先更新全局行情缓存 self.market_cache[snapshot.symbol] snapshot # 再分发给订阅该标的的策略 for strategy in self.strategies.get(snapshot.symbol, []): strategy.on_tick(snapshot) def on_trade(self, trade): # 逐笔成交事件单独分发 for strategy in self.strategies.get(trade.symbol, []): strategy.on_trade(trade)策略实例内部可以维护自己的指标序列、开仓条件、持仓状态等。行情驱动模型的关键在于策略代码必须保证可重入、无状态阻塞——同一个策略实例可能在同一时间被多线程调用因此内部状态要做好锁保护或者干脆让每个策略实例只在一个线程里执行。4.2 信号生成与订单执行策略和柜台之间要隔一层策略算出了买卖信号但这不代表就要立刻下到券商柜台。我把“信号生成”和“订单执行”分成了两层。信号层只负责决定“现在想不想买、想买多少、以什么条件买”执行层负责把信号翻译成订单指令处理撤单、重试、超时、部分成交等事务。为什么必须分开因为实盘执行过程充满不确定性。你拿到的五档行情是时刻变化的信号生成时看到的价格到了下单时刻可能已经完全不同。如果信号层直接把单下出去遇到快速行情很容易买在最高点。更好的做法是信号层发出一个目标持仓或者限价指令执行层根据最新行情动态调整挂单价比如把单子挂在买一价、买二价中间位置或者使用 TWAP/VWAP 算法拆单。另外执行层还要处理“信号重复”问题。一个策略在连续多个行情 tick 里都产生了买入信号如果不做去重执行层可能会连着下好几笔单仓位直接爆掉。解决办法是为信号加触发去重逻辑同一策略、同一标的、同方向的信号在一个有效期内不重复执行。有效期可以基于时间也可以基于持仓状态。4.3 从五档行情里挖掘策略特征几个可以直接落地的思路五档行情接进策略后具体能做哪些特征我分享几个我个人在实盘里验证过的思路。第一个是盘口失衡指标。把买一到买五的总量记为sum_bid卖一到卖五的总量记为sum_ask定义失衡指数为(sum_bid - sum_ask) / (sum_bid sum_ask)。当这个指数快速上升通常说明买方力量在增强可以考虑做多相反则做空或减仓。这里的核心不是指数的绝对值而是变化速度。第二个是主动买盘主动卖盘占比基于逐笔成交里的成交方向。统计最近 N 笔成交中主动买量和主动卖量的比例可以判断当前资金的主导方向。这个数据比只看收盘价涨跌要细致得多因为它反映的是成交时刻的真实买卖力量。第三个是撤单率估计。五档快照是快照你只靠相邻两次快照的持仓变化无法精确还原撤单因为既有新增挂单也有撤单还有成交消耗。但你可以做一个粗估如果买一档的挂单量在前一个快照里是 1000 手价格是 10.00当前快照里价格没变、挂单量变成 300 手同时这段时间内没有对应的成交记录那大概率是有人撤单了。撤单率高往往意味着盘口虚挂严重行情可靠性低。这些特征本身不复杂但把它们工程化、稳定地算出来才是难点。我的经验是先只做最简单的 1-2 个特征跑通全链路再逐步叠加。一上来就堆几十个特征出了问题根本定位不了。5. 工程化实施从原型到稳定运行的七个关键细节5.1 性能观测延迟和吞吐量必须可视化实时行情系统最怕“看起来在跑实际已经不行了”。所以从第一行代码开始就要给系统加上性能观测。我关注的指标有三类。第一是行情延迟。用我在数据模型里提的local_receive_time减去timestamp每 30 秒算一次平均延迟和 P99 延迟。P99 延迟是系统真实的体感延迟平均值容易骗人。比如你平均延迟 20 毫秒但 99% 分位可能有 500 毫秒一旦遇到极端行情就会卡顿这对策略是致命的。第二是队列积压。行情队列的长度和消费耗时是最直观的健康指标。我每 5 秒输出一次队列长度如果持续超过最大容量的 30%就要警惕消费端是不是有瓶颈。等到堆积到 80% 再处理往往已经来不及了。第三是事件处理耗时。把每个策略的on_tick函数耗时记录下来超过阈值就告警。很多时候策略逻辑写的复杂度爆表但没人发现它已经成了整个系统的瓶颈。我的经验是单一策略处理单笔行情的时间不要超过 10 毫秒否则行情一密集就会积压。因此像复杂的统计计算、特征历史序列存数据库等操作都应该异步化或提前预计算不要全都放在实时路径里。5.2 回测和实盘的一致性别让回测骗了你很多玩家在回测里跑出漂亮的曲线一上实盘就变形原因就是回测和实盘的数据流不一致。我踩过最大的坑就是回测用了 K 线实盘却用五档快照策略特征完全对不上信号自然南辕北辙。所以现在我坚持一个原则回测系统必须能消费和实盘完全同一格式的行情数据。做法是给策略引擎做两种运行模式一种是回放模式从历史 tick 数据文件里逐条读取喂给策略引擎另一种是实盘模式从 WebSocket 接收实时数据喂给策略引擎。两种模式的入口统一成同一个函数策略代码完全不用区分自己身处哪个模式。这样回测时是什么行为实盘基本也是什么行为至少排除了数据形态差异这一层误差。历史 tick 数据的获取也要留意。很多免费数据源只能拿分钟线拿不到逐笔成交导致你无法回测真实盘口策略。这时可以先从能拿到的五档快照历史数据开始把回测时间周期调大到秒级先验证策略逻辑是否成立。等到有足够历史 tick 数据了再优化。5.3 运维细节日志、告警、重启恢复一个都不能少实时行情系统不是写完了就能放着不管的它需要稳定的运维支撑。我从一开始就养成了几个习惯你可以直接抄。第一个习惯是结构化日志。每一条日志不能只是一句字符串至少要带上时间、级别、标的代码、事件类型、关键指标值。比如 “2025-02-10 14:30:01.123 | WARN | AAPL | gap_detected | seq_expected8821 seq_received8825”。这样线上排查问题时按时间范围和标的代码一过滤立刻能看到问题节点。第二个习惯是告警分级。行情断连、队列溢出属于高优先级告警必须通过手机通知指标计算超时、延迟偏高属于中优先级记录日志后可延迟处理日常信息类的日志只进文件不需要打扰人。告警阈值不能设得太敏感否则天天半夜被叫醒最后反而忽略真正的故障。第三个习惯是重启恢复。程序崩溃重启后必须能从不完整的市场状态中恢复。最简单可靠的做法是维护一份当前所有标的的最新快照缓存启动时先从缓存加载再开始接收实时行情。这样即使重启策略也能基于最近状态继续运行不会因为缺了 10 秒行情就开始瞎算信号。6. 常见问题与排查技巧实录6.1 典型问题与排查方法速查表结合我自己的项目经历整理了几个出现频率最高的问题和对应的排查方向。问题现象可能原因排查方法行情停止更新但程序还在运行WebSocket 连接断了但没触发重连或者重连逻辑有 bug查看连接的 ping/pong 日志检查是否实现了断线检测测试重连逻辑策略信号明显滞后行情回调和策略计算共用同一线程耗时操作阻塞了接收检查策略处理函数的耗时确认是否用了独立队列消费队列持续堆积数据延迟持续升高消费端吞吐不够或批量数据进来时计算量过大观察队列长度曲线对消费端做性能分析找出慢函数成交量明显偏大或偏小重复数据处理不完全或者数据源有漏发检查去重逻辑是否对逐笔成交生效对比交易所公布的成交量估算回测很好实盘一直亏回测用分钟线实盘用 tick特征不一致或滑点设置过于乐观统一回测和实盘的数据颗粒度在回测中增加合理的滑点与手续费程序偶尔报内存错误队列无限增长或缓存没有清理给队列设置上限检查是否对历史快照无限制缓存多策略同时运行时互相干扰多个策略实例并发修改同一个状态为每个策略隔离上下文加锁保护共用对象或者在事件分发时串行执行6.2 我踩过的三个坑希望你别再踩第一个坑是过度相信行情源的心跳机制。有段时间我用某行情源它的心跳是每 15 秒发一次我天真地以为只要收到心跳连接就没问题。结果某个下午行情源服务端出了故障心跳还在发但行情数据已经不再推送我的策略整整空转了 20 分钟。从那以后我就改成了“数据超时熔断”机制如果某标的时间超过 3 秒没有行情更新立刻告警并暂停该标的的策略交易宁可错过机会也不能用陈旧数据做决策。第二个坑是重连后没有重新订阅。最早的版本里WebSocket 断开重连后只重新鉴权没有重新发送订阅请求导致连接看起来正常但一只股票的行情都不推。这个问题排查了很久才意识到。现在我的重连逻辑固定按“连接 - 鉴权 - 重新订阅 - 校验最后一条行情时间戳”的顺序执行每一步都有日志输出。第三个坑是用 float 存价格导致精度丢失。听起来很基础但真发生在自己身上才肉疼。当时用数据库字段存股票价格显示为 12.3400000001策略判断到价阈值时永远不触发。后来把所有价格字段统一成整数分或 Decimal 类型这类问题才彻底消失。做量化的人一定要养成的习惯凡是涉及钱的数据永远不要用浮点数存储尤其不要做等值比较。6.3 一个小工具行情延迟自检脚本最后分享一个我用得很顺手的小工具。它不依赖复杂框架作用是定时检查某只股票的最新行情是否在推进以及延迟是否在可接受范围。import asyncio import time class MarketDataMonitor: def __init__(self, timeout_sec3): self.timestamps {} self.timeout_sec timeout_sec def update(self, snapshot): self.timestamps[snapshot.symbol] (snapshot.timestamp, snapshot.local_receive_time) def check(self, symbol): if symbol not in self.timestamps: return False, no data biz_ts, local_ts self.timestamps[symbol] now int(time.time() * 1000) if now - local_ts self.timeout_sec * 1000: return False, fstale: {now - local_ts}ms return True, fok, latency{local_ts - biz_ts}ms这套自检逻辑我建议放进主程序的主循环里每 2 秒跑一次。一旦返回异常就立刻走告警通道该暂停暂停该清仓清仓。我自己就是因为这套机制在上线初期拦下了好几次底层数据源故障少亏了不少冤枉钱。行情数据接入这块说实话没有太多玄学核心就是把链路每一环的职责想清楚接收只管接收消费只管消费策略只管算执行只管下单。边界清晰了问题定位、性能扩展都会顺很多。我最初搭系统时也踩了不少坑所以特别理解拿到一个实时行情 API 却不知道从何入手的感觉。希望这篇文章里的思路和代码骨架能让你把自己的组件搭得比我的初版更稳。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进