做实时行情系统这几年,我最大的感受是:它不像普通业务系统那样"能跑就行",而是把延迟、吞吐、可用性三个指标同时按到极致的一类系统。行情晚到一秒钟,交易端的体验就崩了;数据源断掉一次,风控那边就要出大事。这篇文章我就结合自己做过的行情系统,从协议选择、高可用架构到数据源选型,把整套设计思路和踩过的坑完整梳理一遍。无论你是刚接手行情项目的后端工程师,还是准备自建行情服务的团队,这篇文章应该都能给你一份可以直接参考的实操清单。
1. 先想清楚:实时行情系统到底在解决什么问题
1.1 三个核心诉求:低延迟、高吞吐、高可用
行情系统的技术选型,本质上都是在和这三个指标博弈。
低延迟指的是从数据源产生行情到客户端收到行情之间的总耗时。我做过一个交易类项目,业务上给用户的承诺是行情端到端延迟不超过100毫秒,内部设计目标直接压到50毫秒以内。分摊下来,数据源推送到采集层大约10毫秒,采集层处理和入队约10毫秒,推送网关转发约20毫秒,剩余的网络传输和客户端渲染也就剩十几毫秒的空间。每一个环节稍有阻塞,整条链路就超标。
高吞吐更直白。行情高峰期,逐笔成交消息每秒能达到几万条甚至几十万条,同时在线连接数可能是十万级甚至百万级。服务器要反复执行"读消息、解析、路由、推送"的操作,任何一个环节设计成同步阻塞,吞吐就会塌方。实测下来,如果推送网关里混入了数据库操作或者大对象序列化,单机QPS会从几十万掉到几万,完全不是一个量级。
高可用是第三个大坑。行情服务停了哪怕几十秒,用户就会集体投诉;数据丢了哪怕几条,量化策略就可能算出错误信号。我们通常用SLA来衡量,99.99%的可用性意味着一年只能允许约52分钟不可用,听起来挺宽裕,但行情系统的"不可用"往往不是整机宕机,而是某条数据链路上出现短暂抖动。所以设计时不但要保证进程不挂,还要保证数据链路不出现长毛刺。
这三个指标往往是矛盾的:想要延迟最低,就尽量少加中间件,但少了冗余又牺牲可用性;想要高可用,就要多副本、多集群,但数据同步又引入延迟。我的建议是不要在架构上追求单一最优,而是把链路分层,不同层用不同策略,后面会详细展开。
1.2 行情系统的数据分层与业务场景
在设计之前,先把行情数据的类型和消费方盘清楚。
数据层面,行情系统通常要处理两类消息:一类是快照数据,比如某个交易标的的当前价、最高价、最低价、成交量,这类数据可以定期推送,比如每100毫秒或500毫秒推一次;另一类是逐笔数据,也就是每一笔成交的价、量、时间,这类数据是事件流,没有固定频率,行情越活跃频率越高。
消费方也分好几拨。面向普通用户的App,需要的是低延迟快照,用户体验主要看推送间隔和渲染是否卡顿;量化策略系统需要的是逐笔委托和逐笔成交的原始数据,每一笔都可能成为触发信号,所以要求严格有序、不丢不重;风控系统更在意数据完整性,哪怕晚一点,也不能出现缺口。
因为各方的诉求不同,我倾向于在系统内部做一次数据分层。底层统一接入数据源,经过清洗、校验后进入消息总线;往上分两条路,一条做实时快照计算,输出给推送网关;另一条保留全量逐笔流,供量化侧订阅。这样既保证了不同场景的数据质量,又避免了一个消费方把整体链路拖垮。
也可以简单理解成:采集层负责"把数据搞进来",计算层负责"把数据加工好",推送层负责"把数据送出去"。每一层各管各的事,出了问题也能快速定位。
2. 协议选择:没有最好的协议,只有最合适的协议
2.1 主流协议的横向对比
很多刚做行情系统的同学,第一反应是"用WebSocket不就行了吗",实际上协议选择要看数据的源头在哪一端、消费方是谁、网络环境怎样。我把常见的几类协议放在一起对比过,各有各的适用场景。
WebSocket是目前对外推送最主流的方案,浏览器和App都能直接用,全双工通信天然适合服务端主动推送。它基于TCP,顺序有保证,但帧格式有额外开销,连接数高了以后也需要注意内存占用。
TCP私有协议适合源头接入或内部服务间通信。比如交易所或数据源提供的原始行情订阅,很多是基于TCP和自定义二进制协议,优点是协议头可以做到极其精简,延迟低、解析高效,缺点是开发和维护成本高,字段变更需要双方同步修改。
UDP组播常用于证券行情局域网分发。它能把一份行情同时发给多个消费端,网络开销极低,但UDP本身不保证可靠交付,需要自己在应用层做序号和重传补偿。如果网络质量不好,组播丢包会非常头疼。
MQTT是物联网场景下的常用协议,轻量、支持发布订阅,也支持QoS级别。不过行情系统对延迟和吞吐的要求一般高于MQTT的日常应用场景,我用下来觉得它更适合弱网环境或移动端通知,不太适合高吞吐的核心行情链路。
FIX/FAST是金融领域的老牌协议。FIX报文可读性尚可,但字段多、体积大,FAST是针对FIX的压缩编码。这类协议在机构间交易场景很常见,如果你的上游数据源就是FIX协议,对接是躲不掉的。
| 协议 | 延迟 | 吞吐 | 开发成本 | 可靠性 | 典型场景 |
|---|---|---|---|---|---|
| WebSocket | 中低 | 中高 | 低 | TCP可靠 | 对外推送、手机端 |
| TCP私有协议 | 低 | 高 | 高 | TCP可靠 | 上游数据源、内部服务 |
| UDP组播 | 极低 | 极高 | 高 | 弱,需自研补偿 | 局域网行情分发 |
| MQTT | 中 | 中 | 低 | 支持QoS | 移动端、弱网场景 |
| FIX/FAST | 中低 | 中 | 中高 | TCP可靠 | 机构、交易所对接 |
实际项目中我使用的是混合方案,也就是"上游TCP私有协议接入,内部消息总线以二进制消息传递,对外统一走WebSocket推送"。一句话总结:协议没有绝对好坏,完全取决于你手里的数据从哪来、要到哪去。
2.2 对外推送为什么首选 WebSocket
对外推送层,我几乎没有犹豫就选了WebSocket。原因其实很现实:客户端生态太成熟了,浏览器自带WebSocket API,移动端也有成熟的库,团队不需要为每个客户端折腾一套私有协议栈。
WebSocket的全双工能力很关键。行情系统需要服务端主动向客户端推送数据,如果用HTTP轮询,延迟和连接开销都扛不住;如果用TCP长连接私有协议,客户端SDK要自己处理粘包拆包、鉴权握手,工程成本高出一大截。WebSocket把这些问题都标准化了,握手阶段可以带鉴权参数,连接建立后服务端随时可以推送。
另一个优势是WebSocket天然兼容TLS,也就是wss。行情数据虽然不是隐私信息,但加解密能防止数据被中间人恶意篡改,尤其在公网环境下,我建议一律启用wss,不要裸用ws。开了TLS之后,握手开销有所增加,但可以通过连接复用、减少频繁重连来弱化影响。
不过WebSocket也有一个容易被忽视的坑,就是服务端连接数。每个WebSocket连接都是一个持有TCP连接的对象,如果客户端规模有几十万,Gateway的内存占用会非常可观。我们压测过,一个默认2G堆内存的Java网关,单机最多扛几万连接,再往上就频繁GC。后来调整了Netty参数,减少堆外内存拷贝,才把单机连接数提上去。选型时一定要把连接数作为核心容量指标去估算。
2.3 消息格式与增量快照设计
协议确定后,接下来是消息格式。我最早图省事,直接用JSON推送行情。结果压测一跑,序列化开销就把延迟拉高了。行情消息往往高频小体积,JSON的冗余字段和字符串解析成本在这种场景下很吃亏。后来改成了二进制编码,消息体只保留必须字段,解析耗时降了一个数量级。
即便在对外推送层,我也建议优先考虑二进制格式。如果团队实在没有跨语言的序列化方案,可以选择Protobuf或者FlatBuffers。Protobuf生态号,前后端都能直接生成代码解析,缺点是序列化后不够直观,调试时要转成文本。FlatBuffers支持零拷贝读取,对延迟更友好,但接入成本稍高。
消息结构上,我习惯分成基础头、数据体两大部分。基础头包括消息类型、产品代码、发送时间、业务序号;数据体根据消息类型不同,存放快照或逐笔数据。每一类消息都必须带上业务序号,这是断线补数据的生命线。
举个例子,一个快照消息可以设计成这样:
{ "type": "snapshot", "symbol": "AAPL", "seq": 123456, "timestamp": 1699999999999, "data": { "last": 190.25, "bid": 190.24, "ask": 190.26, "volume": 12345678 } }逐笔成交消息则带有独立的tradeId和成交方向。字段设计时,我特别强调不要把所有数字都用字符串表示。像价格、成交量这类高频字段,能用整型就用整型,比如价格放大10000倍后以long型传输,既能保留精度,又能减少转码开销。
增量快照的设计也很关键。全量快照每100毫秒或500毫秒推一次,增量消息紧随其后。客户端先收全量建立初始状态,再按增量逐条更新。一旦发现序号跳变或本地状态对不上,就向服务端请求一次全量快照,把状态拉齐。这是行情推送中最常见、也最稳定的可靠性方案。
2.4 可靠性处理:心跳、超时与断线补偿
WebSocket虽然是TCP连接,但还是会碰到网络闪断、中间设备回收空闲连接、服务端重启等情况。只靠TCP不一定会触发断开,应用层必须做心跳机制。
我常用的心跳方案是:服务端每15到30秒发送一个Ping帧,客户端收到后回Pong;服务端如果在两个Ping周期内没收到任何数据,就判定连接过期,主动关闭并清理会话。客户端侧也要设置超时检测,如果连续多个周期没收到服务端心跳,就主动重连。心跳间隔要结合网络实际情况调整,太短会造成不必要的流量和CPU开销,太长又会让断开检测很迟钝。我们最终线上用的是30秒。
断线重连只是第一步,重连之后的数据补齐才是重点。客户端重连成功后,服务端要从会话存储中取出它最后一次收到的业务序号,然后把后续消息补推给客户端。如果消息总线已经做了持久化,可以按序号范围直接拉取;如果没做持久化,至少要保证能接收一次全量快照,让客户端重新建状态。
消息业务序号的连续性要贯通整条链路。如果业务序号是在采集层生成的,那么后续所有副本都要沿用这个序号,不能内部再生成一套新序号。否则一旦做链路切换,客户端会因为序号混乱而无法判断自己缺了多少数据。这一点的坑,我在后面的常见问题里还会再提。
3. 高可用架构:从单点到多活
3.1 可用性目标与冗余思路
很多小规模行情系统最初都是一个单进程:采集线程拉数据,内存里存最新快照,再用一个WebSocket服务推给客户端。这套方案能撑到几千连接,但必然挂。我接手过一个项目,一次机房网络抖动导致行情停了五分钟,业务方直接炸锅。
高可用设计的核心就一个字:冗余。进程要有备份,数据要有副本,链路要有旁路。但冗余不是简单多开几个实例就完事,还要考虑多实例之间如何选主、如何同步状态、如何切换。我通常按"采集层-计算层-推送层"三个层面分别设计冗余策略,每一层的故障模型不一样。
可用性指标也要提前算清楚。如果业务要求99.99%,那么一年停机时间不能超过约52分钟。这52分钟要分摊到计划内维护、故障切换、数据补齐。所以很多实现细节都要为这个预算让路。比如凌晨的版本发布也算停机,如果需要做到不停机,就必须支持滚动发布和连接优雅迁移。
我的经验是先别急着做跨地域多活,那是成本极高的事。先把同机房内所有单点干掉,做到进程级高可用,再考虑同城双活,最后才是异地容灾。架构上每前进一步,复杂度都成倍增加,不是所有业务都值得。
3.2 采集层的主备与仲裁
采集层是整个系统的入口,它一旦故障,后续全部断粮。连接数据源的采集模块必须做主备。
经典方案是主备两套采集进程同时启动,但同一时间只有主进程在接收并转发数据,备进程处于热备状态,持续检测主进程心跳。当主进程心跳超时,备进程自动接管,重新订阅数据源并开始推送。这里的难点有两个:一是如何判断主进程真的挂了还是只是网络抖动,二是如何避免两个进程同时往外推送数据,俗称"脑裂"。
为了解决脑裂,需要引入一个仲裁机制。小规模团队可以用Redis分布式锁或者ZooKeeper临时节点来做选主。持有锁的节点成为主节点,失锁后自动降级。因为Redis会过期,网络分区时旧主可能短暂仍在工作,所以下游消费端必须对相同序号做幂等去重。
另一种更稳的做法是双采集双写。两台采集同时从数据源接收行情,同时写入消息总线,但各自分配不同的partition范围,下游消费者合并时按业务序号去重。这样即使一台采集彻底宕机,另一台的数据也不会断,切换完全透明。代价是双倍的上游订阅费用和双倍的消息量,数据源配额有限时未必走得通。
我实际用的方案是"主备+心跳+Redis锁",备机平时在待命状态,每秒钟检查一次主进程状态。这个方案实现简单,切换时最多丢几秒数据,配合补数机制可以接受。
3.3 推送层的无状态化设计
推送层是连接数最多、最容易成为瓶颈的一层。庞大的连接状态如果不设计好,想在故障时快速切换几乎不可能。
我的原则是推送网关要做到"逻辑无状态"。所谓逻辑无状态,不是说连接不存在,而是说每个连接的核心信息,比如客户端ID、订阅列表、最后收到序号,都能被外部存储重建。当一个节点宕机后,客户端重连到另一个节点,新节点可以恢复它的订阅关系和补数位点。
但严格把所有连接状态都外置到Redis,每次消息路都查一次,性能会很难看。所以实践中我采用"本地缓存为主、外部存储为辅"的策略:运行中连接状态保存在本地内存,以最快速度推送;同时把关键位点异步上报到外部存储。节点正常时,外部存储只是备份;节点宕机时,用外部存储恢复会话。
为了减小重连后的恢复成本,客户端在重连请求里带上自己最后的seq,新的接入节点可以快速判断:如果本地没有这个连接上下文,就从外部存储或消息总线的缓存队列里拉取后续数据。这种方式比全量重建要快很多,实测切换时间能控制在几秒内。
推送层的容量规划也不能忽略。每增加一万连接,至少要多预留1到2GB内存。如果用的技术栈是Java,还要给GC留出足够空间,否则连接一多JVM就频繁Full GC。
3.4 多活集群与数据一致性设计
到了更大规模,单一机房的吞吐和容灾能力会不够,这时候要考虑多活集群。但多活不是简单地部署两套系统,真正的难点是怎么让两个机房的数据保持连续一致。
行情数据的主链路可以这样设计:两个机房各自部署一套采集和推送服务,同时从数据源接入行情;两套系统产生的消息都进入各自机房的消息总线,再通过跨机房协议将消息流双向同步到对方的主题中。消息带上全局唯一的业务序号,消费者在处理时主要按序号去重,如果发现同一序号出现在两个机房,只处理第一次到达的那份。
网络抖动会导致跨机房同步延迟,所以跨机房链路最好采用数据压缩和批量传输,降低专线占用。如果业务允许,可以设定一个"容忍延迟窗口",比如收到消息后最多等20毫秒再处理,以等待另一侧消息到达,从而减少重复处理。但这个延迟会增加端到端耗时,需要业务方共同评估。
多活最怕的是网络分区。两个机房之间断连,两边都在继续生成行情,下游就会看到两份重叠的流。我处理这类问题的思路是强制角色划分:即使叫多活,在同一时刻仍然只有一个机房承担"主生产"角色,另一个机房处于"热备消费"状态,只在主角色故障时才激活。这样做虽然损失了部分"多活"意义,但换来了最大程度的可控性。
3.5 故障切换不能靠玄学:可观测与演练
高可用不是靠"到时候再处理"撑起来的。没有一套完整的可观测体系,故障发生在哪一层你都说不出来。
监控指标至少要覆盖这些:数据源连接状态、单位时间收到的行情条数、消息序号缺口、采集层到总线延迟、总线消费积压、推送网关连接数、推送延迟、客户端断线重连次数。每一项都要有阈值和告警。
我建议重点盯两个综合指标:一是端到端延迟分布,特别是P99和P999,因为平均延迟永远好看,毛刺都藏在尾延迟里;二是单条链路上业务序号的连续性,只要出现缺口,就意味着数据丢失或乱序,必须立刻告警。
有了监控,还要定期做故障演练。我最常做的演练包括:杀掉主采集进程、模拟数据源断连、拔掉跨机房专线、把推送网关的机器直接重启。演练不只是验证组件能切换,更重要的是验证人的操作手册不依赖当时的记忆。每次演练后都要复盘切换时间、影响范围,以及是否存在"脚本能跑但没人敢执行"的情况。
在故障切换这件事上,我还有一个体会:宁可自动化流程慢一点,也要稳。完全自动切换在极端情况下可能反复抖动,反而把系统搞乱。我们线上采取的是"自动检测、半自动切换"策略:系统监测到异常后自动告警,并生成切换建议;值班人员确认后再触发切换脚本。这样既保留了速度,也给人留了判断空间。
4. 数据源选型:决定整个系统上限的环节
4.1 数据源类型与优劣对比
行情系统里最难替换的往往不是代码,而是数据源。代码写得再差还能重构,数据源一旦不符合需求,整个系统的性能上限就被卡死了。
常见的数据源大概分四类。第一类是官方行情接口,权威性最高,数据质量最有保障,但通常存在配额限制、连接数限制,且不一定能提供细粒度的逐笔数据。第二类是专业数据服务商,能提供全市场、多品种的聚合行情,协议也相对友好,但费用不低,且不同服务商之间的数据质量和延迟差异很大。第三类是第三方聚合接口,接入简单,适合快速验证,但延迟通常偏高,数据细节可能被简化,不适合高精度场景。第四类是自建的行情采集节点,通过合法订阅或公共数据源(必须有合法授权)自行采集并清洗,灵活性最高,但自己承担全部运维成本和合规风险。
无论选哪一种,我都要强调一句:必须确认数据授权范围。不要为了省成本去接入来路不明的数据,也不要超过授权的使用范围把数据二次分发。这个层面一旦出问题,不是技术能解决的。
| 数据源类型 | 延迟 | 数据质量 | 成本 | 接入复杂度 | 适用场景 |
|---|---|---|---|---|---|
| 官方接口 | 低 | 高 | 中高 | 中 | 核心交易场景 |
| 专业服务商 | 中低 | 高 | 高 | 中 | 多市场聚合、机构业务 |
| 第三方聚合 | 中高 | 中 | 低 | 低 | 快速原型、展示页 |
| 自建采集 | 低 | 可控 | 中 | 高 | 深度定制、特殊策略 |
我的建议是,生产环境至少保证一个官方或专业级别的数据源作为主源,再配一个不同来源的备源。两个源不能是同一个上游厂商的同一套接口,否则上游一出事,两个源一起断。
4.2 数据源评估的三个维度
数据源选型不能只看一张宣传页,必须拿实际数据去测。我评估一个数据源时,基本只看三个维度。
第一个维度是延迟。要分别测量接入延迟和更新频率。接入延迟指的是从"源数据产生时刻"到"我们的系统收到数据时刻"的时间差,需要在消息里带上数据源时间戳,用本地时间减一下就能估算。更新频率也很关键,有些服务商宣称实时推送,实际只在秒级聚合后推送,完全达不到逐笔标准。测试时要看行情高峰期每秒能推多少条,以及单条消息的最大间隔。
第二个维度是数据质量。我连续一周按天统计数据源的缺口数、乱序数、异常值数。数据缺口是指业务序号不连续,乱序是后到的消息序号反而小,异常值是指明显突破市场合理范围的价格或成交量。质量差的数据源即使延迟低,也会让下游系统做很多额外的清洗工作,整体算下来反而更贵。
第三个维度是服务质量。需要关注对方的SLA承诺、工单响应时间、定期维护窗口是否频繁、是否有沙箱环境可以测试。我吃过一次亏,某数据源每周三凌晨维护半小时,恰好是我做夜间压测的时间,一测就断流,后来才在文档深处找到维护计划。
这三个维度可以量化为一个选型评分表:延迟占40%,数据质量占40%,服务质量占20%。连续测试一周之后,把每天的最大延迟、平均缺口数、故障次数加权,基本就能筛出靠谱的数据源。
4.3 多数据源冗余与交叉校验
主源和备源选定之后,不能只是"主源挂了人工切换",那样太慢。更稳的做法是让备源在后台持续接收数据,虽然不下发到生产链路,但一直运行着,这样切换时才能保证数据连续性。
交叉校验是双源方案的核心。我日常会校验三组指标:同一标的最新成交价偏差是否超过阈值,通常设置为合理价格区间的千分之一到千分之二;同一时间窗口内的成交量和累计成交笔数是否在合理差异范围内;消息序号是否持续增长,有没有长时间没有新序号的情况。
一旦校验失败,不能立刻切换主备。系统要先判断是自己网络的问题还是数据源的问题。常见做法是连续三次采样失败、且持续时间达到预先设定的阈值,才触发自动切换。切换后还要持续观察备源的数据质量,如果备源数据同样异常,就要停止切换并进入全链路告警。
切换后如何回切也需要设计。我建议不要自动回切,而是由运维人员确认主源恢复稳定后,在低峰期手动执行回切。原因很简单,主源刚恢复时可能还会抖动,自动回切容易造成反复切换,比一直用备源更危险。
数据源这块我还要提醒一个容易被忽略的细节:即使有主备双源,也一定要保留一条直连主源的旁路通道。这条路不参与正常生产,只用于排障。有一次线上数据异常,所有监控都在告警,但完全不知道问题出在数据源还是自己的解析层。我直接启动旁路脚本,在接入层之前单独打印原始报文,一对比就确认是数据源解析方式变了。这个旁路通道关键时刻能救命。
5. 实战落地:一套可运行的行情系统骨架
5.1 核心组件与技术选型
讲完概念,来看一个实际可落地的系统骨架。我以一个中等规模的行情系统为例,目标支撑10万连接、单日处理消息量10亿条左右。
接入采集层:连接数据源的模块,用Go或Java皆可。Go在并发连接和内存控制上更有优势,Java的生态和排查工具更成熟。我们这边用的是Java,基于Netty处理二进制协议,线程模型简单,压测表现稳定。
消息总线层:选用Kafka或Pulsar。Kafka吞吐高、生态好;Pulsar的存算分离和多租户特性在多人共享集群时更友好。我们当时Kafka已经在线运行稳定,所以继续用Kafka。关键配置是把行情主题的分区数设置得足够大,建议至少和下游消费者实例数一致,否则会出现某个消费者热点。
实时计算层:如果涉及快照聚合、多源合并、指标计算,可以用流处理框架。简单场景直接写消费者程序也行,关键处理好幂等和乱序。不建议为了用框架而上框架,行情链路每一跳都会增加延迟。
推送网关层:独立部署一组无状态服务,对外提供WebSocket接入。常见方案是基于Netty或Go原生库做推送网关,单机容量取决于连接数和消息量,10万连接至少准备5到8个节点。
缓存和存储层:最新快照放Redis,命中率高,设置过期时间比如10分钟;历史行情按时序数据库存储,供查询和复盘。如果查询量不大,也可以用ClickHouse做批量导入查询,成本更低。
监控告警层:Prometheus采集指标,Grafana做看板,配合Alertmanager发告警。日志统一走集中式收集,避免故障时逐台机器翻日志。
核心组件清单如下:
- 数据源接入:Netty + 自研协议解析
- 消息总线:Kafka(分区数由下游消费实例决定)
- 实时计算:轻量消费者程序或流处理框架
- 推送网关:Netty WebSocket服务
- 状态存储:Redis(会话备份、最新快照)
- 历史存储:时序数据库或ClickHouse
- 监控:Prometheus + Grafana + Alertmanager
5.2 关键配置与参数参考
很多细节问题不是架构问题,而是参数没调对。我整理了线上稳定运行的一组参考配置,不一定对所有场景通用,但可以作为起点。
TCP参数方面,服务端建议开启TCP_NODELAY,禁用Nagle算法,避免小消息被延迟合并。接收和发送缓冲区可以适当加大,我一般设置SO_RCVBUF为4MB、SO_SNDBUF为4MB,减少网络吞吐瓶颈。Linux内核层面的文件描述符和连接队列参数也要同步调大,具体数值取决于压测结果。
WebSocket心跳间隔设为30秒,服务端在连续60秒内没有收到客户端任何报文时判定超时并断开。客户端如果连续90秒没有收到服务端心跳,主动重连。这样重连频率不会太高,也能及时发现死连接。
Kafka生产端设置acks=all,保证消息写入多副本后才返回,避免单副本故障丢数据。消费端要开启手动提交offset,并确保业务处理成功后再提交。行情场景下,宁可偶发重复消费,也不能丢消息,所以消费幂等由业务序号去重兜底。
Java服务端重点调GC参数。消息量大的时候,对象创建非常频繁,我建议使用G1收集器,并适当增大年轻代和堆内存。更激进的做法是使用堆外内存、对象池减少GC压力。有一版系统在推送高峰期频繁Full GC,我们把Netty消息体改成堆外ByteBuf并启用对象池后,GC停顿明显下降。
监控阈值参考:端到端延迟P99超过100毫秒告警;Kafka消费积压超过几千条告警;业务序号每分钟缺口大于0告警;推送网关单节点CPU使用率超过80%持续5分钟告警。
5.3 端到端延迟优化实战
延迟优化这件事,最怕没有量化目标。我习惯把延迟拆成几段:数据源到接入节点、接入节点到消息总线、消息总线到消费端、消费端到推送网关、推送网关到客户端。每一段都单独探针打点,用日志或者指标记录下来。
优化优先级也很明确。第一优先处理网络传输,减少跨机房跳数、减少不必要的网络转发。如果主链路在同机房,数据源接入和推送网关尽量放在同一可用区,这样RTT可以控制在1毫秒以内。第二优先处理序列化和反序列化,用更紧凑的编码替代JSON。第三优先处理线程模型和阻塞点,杜绝锁竞争、杜绝在IO线程里做耗时操作。
还有一个很容易被忽略的点:对象创建和垃圾回收。如果消息处理过程中创建了大量临时对象,GC会周期性抢占用CPU,导致延迟毛刺。解决方案是使用对象池、复用消息容器、尽量使用基本类型数组代替包装类。
批量处理也能明显降低延迟。消息总线到推送网关之间,如果每条消息都单独消费、单独推送,网络包很小,系统调用开销占比会很高。我们采用批量拉取、批量推送的策略,单批处理几十到几百条消息后再统一网络发送,整体吞吐能提升一个量级。但批量不能无限大,否则单条消息的等待时间会拉长,需要压测找到平衡点。
压测方法上,建议先做单机压测,再做全链路压测。单机压测时重点观察不同并发连接数下的P99延迟和吞吐;全链路压测时重点观察端到端延迟在峰值消息量下是否仍能达标。压测数据一定要接近真实行情,不能只推固定频率的模拟消息。
5.4 上线前检查清单
行情系统上线前的检查,比普通业务系统要多得多。我每一条都吃过亏,列出来供参考。
第一,功能验证:消息字段解析是否与数据源文档一致,快照更新是否及时,逐笔成交是否有序,客户端重连后补数是否准确。第二,性能压测:按预估峰值的1.5到2倍做压测,持续至少30分钟,观察延迟、吞吐、连接数、CPU、内存和GC情况。第三,容灾演练:至少演练主采集宕机、数据源断连、推送网关宕机三个基础场景,确认切换耗时和数据缺口是否符合预期。
第四,监控告警验证:确认告警能正常发出,而不是哑弹;确认值班人员知道如何响应,最好把排查手册和切换脚本放到统一位置。第五,回滚方案:推送网关升级时,老版本是否还能快速恢复;如果数据库结构有变更,是否有兼容旧版本的回滚脚本。
风险控制上,我建议上线前三天做一轮小流量灰度。选择一小部分用户连到新集群,观察延迟和错误率,稳定后再全量切换。行情系统影响面大,宁可多花一天灰度,也不要上线后所有人一起卡顿。
6. 常见问题排查与避坑实录
6.1 连接频繁断开,客户端一直在重连
这个现象很经典。客户端没有主动断,但服务端连接几分钟就消失一次,全网大量重连。排查时先看网关的连接日志,确认是服务端主动关闭还是客户端关闭,还是中间网络设备悄悄掐断。
最常见的原因是服务端心跳超时判断太短。如果客户端网络环境是弱网或者经过移动网络,一个Ping发出去可能要一两秒才能回来,心跳超时设置成5秒就会出现大量误判。我们把服务端判定超时的时间从10秒调整到60秒后,断开率明显下降。
还有一种情况是客户端和服务端之间经过了一些空闲连接回收策略不合理的中间网络设备。连接长时间没有数据就会被认为是空闲连接被清理,即使有WebSocket心跳也可能被忽略。解决方式是缩短心跳间隔到15到20秒,同时客户端针对连接断开做指数退避重连,避免一旦断开,所有客户端同时重连压垮网关。
6.2 行情数据有缺口和乱序
数据缺口通常表现为某个标的的行情突然停更几秒,恢复后价格跳变。乱序则表现为新消息的价格落后于之前已收到的消息。这两者都会严重干扰量化策略。
排查时先定位缺口出现在哪一段链路。在接入节点、消息总线、推送网关分别检查业务序号,看从哪一段开始不连续。如果接入节点收到的序列就缺,基本上是数据源或网络抓包环节的问题;如果接入节点完整但总线消费端缺失,就是消费端处理逻辑或offset提交失误;如果总线完整但推送网关缺,多半是内存缓存淘汰策略把数据丢了。
乱序的产生主要有两个原因:一是多数据源同时接入时,各源的时间戳和序号体系不一致,下游合并时没有按统一业务序号排序;二是消息总线内分区分配不合理,同一标的的行情被路由到了不同分区,消费时无法保证全局有序。解决思路是给每个标的绑定固定分区,同时在消费端做序号校验,发现乱序先缓存等待,超时后再丢弃或补拉。
6.3 多数据源相互打架
接入双源之后,最头疼的问题是主备两个源给的行情不完全一致。价格差几个tick、成交量差一点、快照时间各说各话,交叉校验一直报警。
这不是系统bug,而是不同数据源之间天然存在采样时点和聚合逻辑的差异。有些数据源推送的是最新成交,有些推送的是基于订单簿计算的理论价;有些成交量按笔数统计,有些按股数统计。要解决这个问题,不能只比最终数值,而要明确各自的字段语义,再在做比较前进行口径统一。
对于确实应该一致的核心字段,比如成交价和成交时间,如果偏差超过阈值,我建议按"来得更快且与历史序列更连续"的源作为优先值,另一个源进入告警观察。不要试图在逻辑里动态切换每个字段,很容易把自己绕晕。
6.4 行情高峰期的短暂卡顿
每到行情剧烈波动的时刻,系统就出现几十毫秒甚至几百毫秒的卡顿,但平时完全正常。这种问题往往是某个组件在流量上涨后到达临界点。
最典型的元凶是JVM GC。积累了大量临时对象后,Full GC会暂停所有业务线程,直接表现为推送延迟骤增。排查时看GC日志和JVM监控,如果确认是GC问题,按前面说的调整堆内存、启用对象池、改用堆外内存,通常能得到改善。
另一个元凶是消息总线的消费端处理不过来。行情一波动,积压立刻上涨,消费端如果还在逐条解析和推送,延迟自然越来越高。解决方式是提高批量拉取能力、增加消费者实例、把非核心逻辑异步化。
还有一种情况是网络入向流量打满。一旦某个数据源在峰值时推高频率的大量行情,单机网卡会被占满,其他消息也会被堵住。排查时需要看网卡丢包和软中断占用,必要时给接入层单独划分机器,避免与推送层共用网络资源。
| 问题现象 | 可能原因 | 排查步骤 | 解决方法 |
|---|---|---|---|
| 连接频繁断开 | 心跳超时过短、中间设备回收连接 | 查看双向关闭日志、测试心跳RTT | 调整心跳间隔、客户端退避重连 |
| 数据缺口/乱序 | 链路任一环节丢数据、分区路由不合理 | 分段检查业务序号 | 绑定分区、消费端序号校验 |
| 双源价格不一致 | 字段口径不同、采样时点不同 | 对比原始报文字段语义 | 统一字段口径、按源优先级处理 |
| 高峰期卡顿 | JVM GC、消费积压、网卡打满 | 看GC日志、消费积压、网卡丢包 | 调参、增加实例、网络隔离 |
6.5 数据源连接串流导致的重连风暴
最后一个值得单独说的坑,是数据源偶发断开时,如果没有控制重连频率,会在短时间内把所有采集节点的连接请求同时砸向数据源。数据源本身可能还处于不健康状态,被这一波重连请求一冲,直接拒绝服务,形成重连风暴。
处理方案是给重连加上退避策略。第一次重连等1秒,第二次等2秒,之后按指数退避,上限到60秒。同时把限流逻辑做在采集层,无论数据源多着急,同一分钟内重连次数不能超过阈值。这样做虽然会让恢复时间变长十几秒,但能保证数据源不会因为我们的重连而彻底不可用。
另外,重连成功后不要立刻全量订阅所有标的,可以先订阅一小部分做健康检查,确认数据源恢复正常后再订阅全量。这个动作能防止数据源刚恢复时被瞬间顶到高负载再次宕机。
7. 一些真实体会与建议
行情系统做久了,我有一个很深的体会:技术方案再漂亮,也比不上对真实链路细节的把握。你提前想到了消息切面、双源切换、心跳超时,系统就能多一分稳定;你漏掉一个GC参数、一个序列号校验、一次维护窗口,线上就会在某个深夜教你怎么做人。
如果只能给一条建议,我建议先把"可观测性"放到最高优先级。没有完整的指标和日志,一切高可用设计都是盲人摸象。我见过太多团队,架构图画得很完整,线上出问题时却连数据在哪个环节丢了都不知道。
最后再分享一个小技巧:在行情系统正式上线后,保留一个从数据源直连到独立脚本的旁路通道,不参与生产链路,只用于排障和对比。这个通道看起来浪费资源,但在你面对异常数据、怀疑某个环节出了问题时,它是最高效的定位工具。我靠这条路解决过至少三次疑难杂症。
行情系统是一个持续演进的过程,不用急着一步到位。先把数据接进来、推出去,再把可靠性做扎实,最后根据业务需要逐步完善多源和多活能力。每一步走稳了,系统自然会越来越强。