十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

实盘监控行情接入:从单股API轮询到批量实时推送的工程实践

实盘监控行情接入:从单股API轮询到批量实时推送的工程实践 做量化实盘监控的第一版我图省事直接循环去调单只股票的行情接口。10只股票没问题20只勉强能跑到50只的时候系统彻底不对劲了——行情刷新越来越慢动不动就报限频错误更坑的是不同股票的行情数据时间点根本不齐排序和筛选用起来心里发虚。后来我把整个架构推倒重来换成了批量订阅推送的方案才算真正把实盘监控跑稳了。这篇文章就聊聊为什么实盘监控不能只调用单只股票行情 API以及批量实时行情在工程上具体应该怎么做。标题里提到的实盘监控、API、批量实时行情、量化工程这几个词本质上是一条链路量化系统需要一个持续、稳定、低延迟的行情源而行情源的接入方式直接决定了监控系统的实时性和可靠性。如果你也正在搭实盘监控或者已经踩了轮询接口的坑这篇文章应该能给你一些有价值的参考。1. 单股轮询在实盘监控里的真实困境不是慢是失真很多刚接触量化的朋友会觉得调用单只股票行情 API 是最简单的方案写个 for 循环把股票池遍历一遍拿到数据就完事了。这个方案在回测、离线分析、低频策略里确实够用但放到实盘监控里问题会一层一层地暴露出来。1.1 轮询周期、接口延迟与数据时钟的错位先想一个问题你看到的行情到底是哪个时刻的行情单股轮询的逻辑是一次请求拿一只股票的最新快照。网络请求发出去行情服务端收到后返回当前最新价这个过程有网络往返时间也有服务端处理时间。假设单次请求的耗时是 150 毫秒你循环请求 100 只股票不考虑并发跑完一整轮就需要 15 秒。也就是说你界面上显示的实时行情第一只股票和最后一只股票之间实际差了十几秒。这还不是最严重的。真正麻烦的是数据时钟的错位。你在计算板块涨幅排名、资金流、异动监控的时候会默认所有股票的行情是同一时刻的快照然后把它们放一起对比排序。实际上第一只股票的数据是 0 秒的最后只是 15 秒后的这两只股票在涨幅榜上的位置对比用一个已经过期的价格去和一个更新的价格比排序结果完全可能失真。我记得自己第一次发现这个问题是在做异动监控的时候。系统提示某只股票一分钟内涨了 5%我点开明细一看最新价还是几分钟前的因为轮询到它的时候刚好卡了一次网络重试。那一刻我就明白了轮询模式下的实时本质上是一个概率事件不是必然结果。1.2 限频的数学账100只股票每3秒刷一次的后果行情 API 普遍有访问频率限制。以我接触过的多数行情源为例免费或低成本的接口通常限制在每分钟几百次请求好一点的行情终端能到每分钟几千次但价格也跟着上去了。我们算一笔账。假设你监控 100 只股票希望每 3 秒刷新一次行情。一秒钟需要的请求次数是 100 / 3 ≈ 33 次一分钟就是 2000 次请求。绝大多数普通行情接口根本扛不住这个量级。那能不能把轮询间隔拉长到 10 秒可以但 10 秒的刷新频率对实盘监控来说已经不太够用了。尤其做短线异动监控10 秒前的价格和当前价格可能差出好几个档位等系统反应过来行情早就走完了。而且即使你把频率降下来遇到开盘、收盘这种行情密集的时段接口响应变慢请求排队实际轮询周期会被进一步拉长限频和超时基本是躲不掉的。并发轮询能解决吗能解决一部分但会引入新的问题。用线程池并发请求 100 只股票确实能把一轮耗时压到几百毫秒但并发请求数一上去更容易触发接口的并发连接数限制。我实测过一个只允许 5 个并发连接的行情接口你用 20 个线程去轮询大量请求会直接失败或者超时整体成功率反而比串行更低。1.3 更大的隐患把监控系统的可用性和外部 API 绑在一起单股轮询模式下每只股票的行情都依赖一次独立的 HTTP 请求。任何一次网络抖动、DNS 解析慢、服务端限流都会直接影响某只股票的行情更新。而你的监控逻辑没办法区分这只股票真的没波动和这只股票的行情没有取到于是出现两种典型的误报一是漏报。股票已经快速拉升但行情请求失败或超时监控系统还拿着旧价格做判断自然发现不了异动。二是误报。股票的行情一直没更新某个字段触发了告警规则你打开一看价格其实是几分钟前的白白浪费注意力。我在生产环境里还踩过一个更隐蔽的坑外部 API 的可用性波动会被直接传导到监控系统的稳定性上。行情源某次升级或者限流策略调整导致部分请求持续超时我的监控系统因为等待响应的线程堆积内存飙升最后把整个服务拖挂了。明明只是行情源的问题却把自己系统的可用性也赔了进去。这也是我后来坚决不用单股轮询做实盘监控的根本原因。2. 批量实时行情的路线选择推送优先轮询兜底既然单股轮询不行那批量实时行情该怎么做我自己的结论是推送优先轮询兜底。核心思路是让行情数据主动来找你而不是你去一遍一遍地取。2.1 三种数据获取方式的核心差异先梳理一下当前主流的行情接入方式对比一下各自的适用场景。对比项单股 HTTP 轮询HTTP 批量轮询WebSocket 推送请求次数N 次/周期1 次/周期1 次连接持续推送数据时效性高延迟受循环顺序影响中等批量接口也有处理延迟低延迟服务端主动推送限频风险极高中等极低断线恢复难度简单重新请求即可简单重新请求即可复杂需要快照增量机制开发成本低低中高典型场景低频策略、离线分析中低频监控、分钟级策略实盘监控、tick 级策略单股轮询的劣势我们上面已经聊过了。HTTP 批量轮询是一次请求同时返回多只股票的行情请求次数从 N 次降到了 1 次限频压力瞬间小了很多。我试过把 100 只股票的代码拼成一个参数一次批量请求平均耗时在 200 到 400 毫秒每 3 秒刷一次完全可行。如果你的监控标的不多、对实时性要求不是极端苛刻批量轮询是一个性价比很高的方案。但批量轮询依然有它的天花板。首先批量接口的返回频率主动权在客户端你始终是周期性拉取两次请求之间出现的价格波动你是感知不到的。其次批量接口的数据也是某个时间点的快照日内高频交易的监控需求它满足不了。WebSocket 推送则是完全不同的模式。客户端和服务端建立一个长连接服务端有行情变化就主动推给你。你不需要管什么时候去取数据只需要处理数据来了怎么用。这种模式天然适合实盘监控延迟低、频率高、没有轮询的周期空档。2.2 WebSocket 推送方案从我去取变成它送来我目前的实盘监控主链路就是 WebSocket 推送。接入的过程有几个关键点值得细说。连接协议方面大多数行情 WebSocket 服务端会遵循类似的交互流程客户端发送订阅请求服务端先返回一批快照数据之后持续推送增量行情。订阅请求一般长这样{ action: subscribe, symbols: [600519, 000001, 300750], channels: [quote, trade] }收到订阅确认之后你就能在回调里不断收到 tick 数据了。行情数据的格式每家不一样但核心字段大同小异我比较关注的是这几类快照类字段最新价、开盘价、最高价、最低价、累计成交量、累计成交额、买卖五档逐笔成交字段成交价、成交量、成交时间、成交方向标识字段交易所时间戳、本地接收时间戳、序列号这里面最关键的是交易所时间戳和序列号这两个字段。做实时监控不能只看数据到没到还要知道这个数据是哪一刻产生的以及有没有跳过了某些中间行情。这两个字段在做数据对齐和乱序处理时特别重要后面我会专门讲。WebSocket 的开发成本比 HTTP 轮询高主要体现在连接管理、心跳保活、断线重连、数据完整性校验这些方面。但一旦把这些基础设施写好后续扩展标的就非常简单了直接往订阅列表里加代码就行不需要为每只股票单独发请求。2.3 多数据源混合架构主推备用轮询怎么切换做实盘监控我强烈不建议只依赖一个行情源再稳定的服务商也有出问题的时候。我的架构是双数据源主源用 WebSocket 推送备用源用 HTTP 批量轮询。主源正常情况下提供所有实时行情备用源保持低频轮询比如每 5 秒拉一次全量快照用来做交叉校验和兜底。两个源的数据会同时写入一个统一的数据结构我可以通过比较两个源的价格偏差来判断数据是否异常。切换逻辑我是这样设计的当主源连续 3 秒没有收到任何行情推送或者心跳超时就认为主源异常自动切换到备用源。切换过程对上层策略透明——策略层只从本地缓存读数据不关心数据到底是哪个源来的。主源恢复后再等它推送的数据连续稳定 30 秒才允许切回。这里有一个特别需要注意的点切换不能光看有没有数据还要看数据新不新。有的场景是连接没断但服务端推送频率突然变得很低比如一分钟才推一次。这种假活状态比连接断开更难发现。我的处理方式是监控每条 tick 的时间戳如果行情时间戳和当前系统时间的差距持续超过阈值比如 5 秒就触发降级切换。3. 行情网关的四个核心模块订阅、缓存、对齐、回放批量实时行情的接入不是写一个 WebSocket 客户端连上就完事了。要支撑实盘监控行情数据在进入策略逻辑之前需要经过一层行情网关来处理。我把这层网关拆成了四个核心模块。3.1 订阅管理动态增删标的不重启服务实盘监控里股票池不是一成不变的。今天关注的 100 只股票明天可能要换掉 20 只。如果每次调整都要改代码、重启服务那这个系统的运维成本就太高了。订阅管理模块要做的事情就是动态维护一个标的集合。外部通过接口向网关发起订阅或退订请求网关负责更新订阅列表并同步给 WebSocket 连接。一个简单的实现思路是这样的class SubscriptionManager: def __init__(self): self._subscribers set() def subscribe(self, symbols): self._subscribers.update(symbols) self._push_subscribe_request(symbols) def unsubscribe(self, symbols): self._subscribers.difference_update(symbols) self._push_unsubscribe_request(symbols) def snapshot(self): return list(self._subscribers)有了这个模块就可以通过管理接口动态调整监控范围不用重启服务也不用担心旧的订阅还占着连接资源。3.2 本地快照缓存让策略层永远只读内存行情网关的核心价值之一是隔离。策略层、UI 层、告警模块不应该直接去解析 WebSocket 原始消息而是应该从一个统一的数据缓存中心读取数据。我把这个缓存设计成了快照增量的双层结构。快照层是一个字典结构key 是股票代码value 是最新的行情快照包括最新价、成交量、买卖五档等信息。增量层是一个队列保存最近的 N 条 tick 数据供策略层做事件驱动分析比如逐笔成交监控、价格突变提醒。class QuoteCache: def __init__(self): self._snapshots {} self._recent_ticks deque(maxlen10000) def update_quote(self, symbol, quote): self._snapshots[symbol] quote def append_tick(self, tick): self._recent_ticks.append(tick) def get_quote(self, symbol): return self._snapshots.get(symbol) def recent_ticks(self, count100): return list(self._recent_ticks)[-count:]策略层永远只读这个缓存不关心行情是从 WebSocket 来的还是从 HTTP 批量轮询来的也不关心行情源是否发生过切换。这种解耦让整个系统的稳定性和可维护性都提升了一个量级。3.3 时间戳与序列号数据对齐的两种手段批量实时行情接入之后最容易让人迷惑的问题就是数据对齐。我给你举个例子某只股票的 WebSocket 推送里有一条成交记录本地接收时间是 14:30:01.200但这条记录在交易所实际发生的时间是 14:30:01.050。如果你用本地接收时间去做分析会引入网络传输的延迟误差。我的做法是统一使用交易所时间戳作为主时间轴。所有本地分析、排序、行情快照的计时都基于 tick 里的交易所时间戳而不是本地接收时间。本地接收时间只用来做延迟监控——计算交易所时间和本地时间的差距用于评估行情源的实时性。序列号则是用来发现数据是否丢失或乱序的。服务端推送的每一条 tick 都带一个递增的序列号客户端按序列号递增的顺序处理。如果发现序列号跳变说明中间有数据丢失需要触发补数逻辑。如果发现后到的数据序列号比已处理的小说明发生了乱序要根据业务场景决定是丢弃还是重新排序。3.4 断线重连后的回放链路全量快照增量补位WebSocket 推送最麻烦的地方在于断线重连。TCP 连接断开期间行情一直在发生重连之后如果不做数据补偿中间这段时间的行情就永远丢失了。对实盘监控来说这可能是致命的。我用的方案是全量快照增量补位。重连成功后网关先发送一个全量订阅请求服务端会返回所有订阅标的最新快照。这个快照能保证行情数据有一个新的基准点。然后再从断开时间点附近拉取增量行情把重连窗口的数据补回来。不同行情服务商对增量补位的支持不一样。有的支持按时间范围拉取历史 tick有的只支持拉取最近 N 条。如果不支持增量补位我就用备用源的批量轮询数据来覆盖重连窗口虽然粒度粗一点但能保证关键价格数据不出现长时间空白。这块的逻辑一定要做好状态管理。我曾经犯过一个大意的错误重连之后直接接收增量推送没有先拉全量快照结果本地缓存还留着断线前的旧价格增量推送又没有覆盖到所有字段最后导致部分股票的监控数据新旧混杂查了整整一个下午才发现问题。4. 实测数据与关键参数从 demo 到能扛住实盘讲完架构聊点实际的东西。我自己在测试环境和生产环境跑下来的一些数据和参数给大家做个参考。不同的行情源、不同的网络环境会有差异但量级应该差不多。4.1 三种方案在100只标的下延迟与开销实测我在同一台服务器上用同一批 100 只股票对三种接入方式做了对比测试。服务器配置是 4 核 8G 内存网络环境是普通的机房带宽。单股 HTTP 轮询每 3 秒循环一次串行执行一轮完整请求平均耗时在 2.5 到 3.5 秒之间。如果你真的这样配置实际刷新周期会变成 5 到 6 秒因为请求耗时和轮询间隔是叠加的。期间还会时不时遇到单只股票请求超时超时重试又进一步拖慢整轮循环。HTTP 批量轮询100 只股票拼成一次请求平均耗时在 300 到 500 毫秒之间。每 3 秒刷一次可以稳定做到延迟主要看批量接口自身的处理速度。这个方案做分钟级监控、中低频信号提醒是够用的。WebSocket 推送从行情在交易所产生到我的服务收到数据延迟中位数在 80 到 150 毫秒左右P95 在 300 毫秒上下。内存开销方面100 只股票的行情快照加 1 万条 tick 的增量队列大概占用 20 到 30 MB 左右完全在可接受范围内。接入方式100只股票实际刷新周期端到端延迟典型限频风险单股 HTTP 轮询5-6 秒含请求耗时秒级高HTTP 批量轮询3 秒可控亚秒级中WebSocket 推送事件驱动无周期100ms 左右低4.2 心跳、超时、重连参数的合理设置WebSocket 连接的稳定性很大程度上取决于心跳和超时参数设置得合不合理。参数太保守连接断了半天发现不了参数太激进又会频繁误判断开造成不必要的重连。我自己用的参数是这样的每 15 秒发送一次 Ping 消息如果在 30 秒内没有收到 Pong 响应或者任何行情数据就判定连接已断开触发重连。这个 30 秒的判断窗口既能覆盖网络抖动又能比较快速地对真实断线做出反应。重连的策略我建议用指数退避最大上限的组合。第一次重连等待 1 秒第二次等 2 秒第三次等 4 秒以此类推但最大等待时间不要超过 60 秒。这种策略的好处是短暂网络抖动时可以快速恢复长时间断线时不会用高频率的重连请求去冲击行情服务器。有一个容易被忽略的点重连成功之后不要立刻开始接收增量行情要等全量快照拉取完成、本地缓存更新完毕之后再恢复策略层的读取。否则策略层可能在快照还没更新的窗口期读到旧数据。4.3 数据源异常识别与自动降级触发条件双数据源架构要想真正发挥作用必须有明确的异常识别和降级触发条件。否则异常发生了系统还在傻傻地用异常数据做判断那备用源就白设了。我总结的触发条件分为三类第一类连接层异常。WebSocket 连接断开、心跳超时、连续多次重连失败。这类异常最直接一旦触发就立即降级。第二类数据层异常。连续 N 条 tick 的交易所时间戳比本地时间旧超过 5 秒或者序列号跳变超过某个阈值比如连续 10 次重连都出现数据回补失败。这类异常说明连接可能还活着但数据链路已经不可靠了。第三类数据源交叉校验失败。主源和备用源在同一时刻的价格偏差超过阈值比如超过 0.5%持续多次。这种情况可能意味着主源价格有误也可能意味着备用源已经过期需要进一步排查。降级动作执行后要输出一条结构化的告警日志记录降级时间、原因、当前使用的数据源、涉及的股票范围。这样出了问题排查的时候有据可查不用靠猜。5. 实盘监控里最容易翻车的三个细节架构和参数都到位了但真正跑实盘的时候还是有一些细节容易翻车。这些坑如果不提前规避关键时刻会给你来一下狠的。我把自己的经历写出来你们能避则避。5.1 行情乱序导致最新价倒退第一次遇到最新价倒退的问题时我一度以为是行情源数据错了。某只股票明明在快速拉升系统里的最新价却突然比前一条低了一大截然后过一会儿又跳回来。后来排查发现问题出在网络传输的乱序上。TCP 协议本身能保证字节流的有序性但 WebSocket 服务端的推送链路里如果经过消息队列、多节点转发不同消息到达客户端的顺序可能和发送顺序不一致。比如一条更新的 tick 先到并更新了最新价紧接着一条更早的 tick 后到如果不做任何判断就直接覆盖最新价价格就会倒退。解决方案就是在更新快照之前检查序列号和时间戳。如果新到的消息序列号小于当前快照的序列号说明是过期消息直接丢弃如果序列号相等但时间戳更旧也是同样处理。只有序列号递增、时间戳更新的消息才允许写入快照缓存。def on_tick(self, tick): symbol tick[symbol] seq tick[seq] current self._snapshots.get(symbol) if current and seq current.get(seq, 0): return # 过期消息直接丢弃 self._snapshots[symbol] tick5.2 在行情回调线程里做 IO 操作把推送线程堵死这个坑我印象太深了。最初版本里我在行情回调里直接写了日志落盘和数据库更新操作想着顺便把数据存下来。结果某天行情剧烈波动一秒推过来几百条 tick每条回调里都要写一次 SQLite磁盘 IO 排队推送线程被堵住了后面的行情越积越多系统延迟从几百毫秒飙到了几十秒直接把实盘监控变成了事后复盘。WebSocket 客户端解析消息、派发回调的线程是非常宝贵的资源。回调函数里绝对不能有阻塞操作包括数据库写入、远程调用、大文件写入。正确做法是回调里只做两件事更新内存缓存、把 tick 放入一个线程安全的队列。真正的持久化操作由单独的消费线程从队列里取数据异步执行。我当时的改造方案是用一个queue.Queue做缓冲回调函数只负责put一个后台线程负责get后批量写入数据库。改造之后行情推送线程的耗时从原来的几十毫秒降到了微秒级系统延迟恢复了正常。5.3 用延迟监控判断监控系统本身是否健康最后一个建议跟业务逻辑无关但特别值得做给行情接入加一套延迟监控。延迟的定义很简单就是本地时间减去tick 里的交易所时间戳。我每天开盘后都会盯这个指标的走势图。正常情况下WebSocket 推送的延迟应该在 200 毫秒以内上下波动。如果延迟突然持续走高说明行情链路出现了拥堵即使连接没有断开系统的实时性也已经打了折扣。延迟监控的价值在于它能让你在数据彻底断掉之前提前感知到链路质量的下滑。我做了一个简单的告警规则连续 30 秒内行情延迟 P95 超过 2 秒触发警告。这个规则曾经救过我一次某天行情源服务端负载异常延迟逐渐走高我在数据完全中断前就切到了备用源避开了那段时间的监控盲区。另外延迟监控的历史数据对复盘也很有用。出了数据问题回看延迟曲线能帮助判断问题是从什么时候开始的、影响面有多大。这比事后翻日志高效太多了。实盘监控的行情接入从本质上看不是一个调接口的问题而是一个需要把数据链路、状态管理、异常处理都考虑进去的工程问题。单只股票 API 在特定场景下自有它的价值但把它直接搬来做批量实时监控最终一定会在实时性、稳定性和可维护性上付出代价。希望这篇文章能帮你少走一些我走过的弯路把更多精力放在真正重要的策略和信号上。
返回列表