
做实时行情系统尤其是给交易、风控、量化策略提供数据支撑的那种有个很扎心的现实行情这玩意儿不像普通接口那样“能通就行”。它对延迟、完整性、可用性的要求几乎可以按军工级来理解。早几年我接手过一个行情网关项目从协议选型到多源校准再到跨机房高可用踩了一路的坑。今天把这套设计思路完整拆出来希望能帮正在做同类系统的朋友少走几步冤枉路。先说清楚这个系统到底解决什么问题。行情源本身不太可能只有一个交易所直连、第三方聚合源、商业数据商可能同时存在每条链路的延迟特征、字段标准、异常行为都不一样。系统要做的就是把不同的数据源接进来统一清洗、校准、缓存再通过一套稳定低延迟的接口往上层推。难点不只在接入更在于多个源打架的时候到底信谁以及某个源挂了以后整体怎么降级而不产生灾难性后果。适合谁来参考正在做行情网关、金融数据中台、量化回测数据管道或者任何强依赖实时数据的后端系统的人这篇都能用得上。1. 协议选择行情数据该走哪条“路”1.1 行情传输的特殊性普通业务接口可以容忍几秒钟延迟行情系统不行。以国内商品期货的高频快照为例一个活跃合约单日可以产生几十万笔tick数据一天的行情文件解压后动辄几个GB这还没算盘口深度数据。更重要的是行情数据天然是时间序列流一旦断流超过几百毫秒行情快照就不连续了上层策略做出来的判断可能整套崩掉。所以传输协议的选择要考虑的不只是“能不能把数据送过去”还包括在极端网络抖动下的恢复能力、丢包补偿能力、以及流式处理下的内存开销。顶层的设计思路我一直坚持“协议分为层”的原则传输层管字节怎么走应用层管消息怎么解析两层不要混在一个自定义协议里。很多团队图省事直接拿TCP长连接传JSON短期内够用一旦并发上来和源数量增多解析性能和流量成本都会成为瓶颈。下面分传输层和应用层两层来说。1.2 传输层协议怎么选TCP、UDP、组播这三类行情行业里是最常见的。TCP胜在可靠性高、内网传输稳定但缺点也很明显队头阻塞。一个包丢了后续所有包都得等重传在行情流这种连续高频的数据场景下哪怕只是偶发的丢包都可能造成明显的延迟毛刺。UDP则相反没有重传机制天然适合实时流但需要业务层自己设计丢失检测和补偿逻辑。组播更特殊一点它只在局域网环境有优势一台机器发多台机器同时收网卡层面就做了分发放大倍数特别高。但组播对交换机、网卡和驱动的要求比较高跨公网基本用不了。我实测下来比较稳的组合是这样的内网核心链路用单向UDP加应用层序列号检测外网分发走TCP长连协议体内带最后一条序列号回执。UDP部分负责低延迟推送TCP部分负责补偿拉取和握手连接两条通道互不干扰。对于交易系统内部的组播场景如果交换机支持IGMP Snooping放行情快照效果非常理想但记得必须配IGMP Querier否则组播流量在某些交换机下会断。这些设计要在初期就定下来不然后期改协议上层的所有客户端、服务端解析逻辑全部要重来。1.3 应用层协议与消息体设计应用层协议是个容易被低估的工作点。很多系统直接用现成的JSON图一时方便但实际上行情数据有“结构复杂、字段固定、重复度高”三个特点。结构复杂意味着JSON序列化和反序列化的CPU开销很高字段固定意味着绝大概率不需要动态增删字段重复度高则非常适合字段剪裁和批量压缩。所以我个人通常会推荐二进制协议最稳的路子是Protobuf或者FlatBuffers。Protobuf胜在成熟、跨语言支持好并且能通过proto文件严格约束schema新增源或者修改字段时不容易出错。FlatBuffers不需要反序列化就能直接读取特定字段适合对读取延迟极度敏感的C/Rust层但接入成本会高一些。如果团队规模不大就老老实实用Protobuf别自己去写二进制序列化手写的编解码器最容易出字节序和位域这种低级错误。协议设计上有几个关键点必须做好。第一消息体里面必须有全局递增序列号这个序列号是上层对账和丢包检测的基础别偷懒。第二必须有明确的消息类型和版本号哪怕是同一个源不同时期返回的消息字段也可能不同没有版本号将来会酿成大祸。第三合理定义消息边界。TCP长连是流式传输接收端需要知道一条行情消息从哪开始到哪结束这就在协议头里必须有长度字段。第四时间戳精度用微秒级不要用字符串用int64存UTC微秒比较好能省带宽也方便排序比较。1.4 不同场景下的协议组合建议协议不是越复杂越好核心是匹配场景。做一套协议组合方案的时候我喜欢先把数据流向搞清楚不同节点的需求是不同的。比如交易所侧接入、内部总线、对外API这几个环节各有各的职责强行统一反而麻烦。我给一个比较通用的组合方案表实际可以按你们自己的情况调整场景传输层应用层备注内网核心行情分发UDP 应用层确认/补偿Protobuf二进制低延迟为主序列号补漏跨机房主备同步TCP长连接Protobuf 压缩快照保证可靠性优先对外API/盘口推送WebSocketJSON或MessagePack兼容浏览器和客户端前端图表推送WebSocketJSON消息频率要限频防浏览器渲染卡顿WebSocket是很多初建团队的首选因为浏览器友好、思维模型是连上就推不会有TCP粘包的问题。但注意行情高频推送时WebSocket的应用层帧开销比较大一条小消息的封帧、掩码处理都要消耗CPU单节点支撑的客户端数量远低于裸TCP。如果技术栈前后端分离并且客户端数量在几千这个量级WebSocket没问题如果你要做的是机构级网关客户端上万那还是考虑TCP私有协议或者MQTT broker做落地推送更好。MQTT在这种场景下有个额外优点就是QoS、遗嘱消息这些都是现成的协议层能力弱网下比裸TCP抗造。2. 高可用架构从单点服务到可容灾的行情链路2.1 行情链路的基本组成先画一个最简的行情链路模型数据源接入层、消息缓冲层、行情处理层、对外分发层、存储层。接入层负责与不同的行情源建立连接做心跳检测、断线重连、原始数据校验。消息缓冲层负责承接数据流量高峰并解耦生产者和消费者Kafka在多数情况下是首选但要注意topic分区数的设置理想情况下每个源或者每组契约独立分区避免因单个分区流量过大拖累整体。行情处理层负责清洗、归一化、事件生成、多源校准是业务逻辑最重的部分。对外分发层是面对客户端的一层需要支持多协议接入同时承担连接鉴权、权限控制、流量控制。存储层则负责把行情落盘供后续回放、回测和分析用注意这里必须支持批量写入。你要是后端服务经验比较丰富会发现这个链路跟普通消息系统很像。但关键不一样的地方在时延预算。普通系统可以在缓冲层等待批量积压后统一处理行情系统不行。从数据源到达行情处理层开始到对外分发层将同一帧行情发出这个端到端延迟通常要控制在几十毫秒甚至个位数毫秒。如果中间引入一层Kafka延迟就会增加不少所以很多极低延迟场景宁可用共享内存、内存队列甚至ZeroMQ这类工具而不是重量级的Kafka。我的经验是业务实时处理通道和持久化通道要分开处理通道走内存级管道持久化通道再做异步批量写不然一次行情洪峰就能把消费者拖死。2.2 关键组件的高可用设计单机跑任何一个组件都是安全事故行情系统尤其如此。先说接入层所有数据源连接必须做成主备双链路主数据源掉线后备份数据源要能在秒级接管状态。这里容易忽略一点备份链路不能只是建立连接后就空跑必须同时接收数据只是对外标记为备份流即可不然主备切换后会有数据空洞。处理层设计需要考虑进程级失败不能用单进程硬抗用多进程双活加负载分担是比较常规的做法。比如同一份数据由两个进程实例同时跑对外暴露同一个虚拟IP客户端连接在主实例上如果主实例挂了VIP漂移到备实例对客户端来说只是一个瞬断。Kafka层面的高可用副本因子建议至少等于3且最小同步副本配置为2这样才能保证在单个broker掉线的情况下读写不中断。但Kafka的高可用只是消息不丢不自动解决消费逻辑的幂等和去重。消费端拿到消息后必须依据行情序列号做去重和乱序重排否则主备切换之后同样的tick被消费两次会导致指标重复计算盘口快照状态也会错乱。这个是所有从单机版演进到集群版最容易踩的坑。存储层高可用相对好做数据库和文件系统都有成熟方案但行情系统最忌讳的是为存储牺牲写入性能。建议时序数据库选型优先考虑InfluxDB、TimescaleDB或者ClickHouse写入模型天然适配高并发时序流。注意存储节点的归档策略有些数据的价值随时间递减可以设置保留周期减少空间压力同时也能避免查询越来越慢。2.3 数据一致性和故障恢复行情系统和数据源之间的状态同步要考虑断连时的处理逻辑。数据源端到端有一个“当前快照”概念而网络闪断期间应用层可能收到了新版本的数据但物理层并没感知到丢包这种情况如果直接丢弃新数据会造成状态不一致。一个常用的方案是快照加增量结合的同步模式恢复连接后先拉取一个全量快照再对齐序列号之后进入增量同步。而这个中涉及到的对齐点不能简单用品最后一帧的时间去对齐因为各源之间时钟本来就存在偏差最好使用数据源自带的序号字段作为对齐基准。还有状态机设计。每个数据通道内部都应有自己的连接状态机例如DISCONNECTED、CONNECTING、SYNCING、READY四个状态。连接断开时状态机切换到DISCONNECTED进行指数退避重连SYNCING阶段说明正在做序列号对齐和快同步READY才对外放量。对外分发层应只在READY状态下向业务网关提供数据这是避免脏数据的核心手段。很多团队成员误以为只要TCP连接没断就代表数据是正确的实际上TCP长连接上多了一层“逻辑连接”必须用业务层心跳和序列号校验逻辑连接的活性。2.4 容量估算与性能参数设计阶段就定性能目标永远不要等活动上线了再看数据。我习惯先把99百分位延迟、99.99百分位延迟、每秒事件处理能力TPS、最大连接数、接收吞吐量这些指标写进设计文档后续所有架构决策都围绕它们来验证。一个简单估算公式假设你有10个数据源通道每个通道每秒能产生5万条行情tick那么系统整体需要处理的峰值事件量是10×5万50万条/秒。如果每条行情是200字节的序列化消息每秒输入带宽大约100MB处理层消费这份流量时内存和GC损耗可以按3倍估算。如果采用Java技术栈这种情况下至少需要16个CPU核的节点、32GB堆内存起步且GC必须选用低延迟垃圾回收器比如ZGC。Kafka分区数可以按总吞吐量除以单个分区安全吞吐量估算通常一个分区的单生产者写入吞吐量会稳定在每秒几千到几万条那么至少需要16到32个分区才能安全承载这每秒50万的事件量。算完之后再看线程池、客户端连接数、序列化吞吐量这些细项逐步逼近一个安全值。千万不要用感觉拍脑袋定配置行情系统的性能瓶颈一定是算出来的。3. 数据源选型与多源校准数据准不准比快不快更头疼3.1 为什么非得搞多数据源单一数据源最大的问题是“你没法验证你拿到的数据是对的”。行情源偶尔会推错数据比如某一次源库切换实例触发了一条系统重启后的首帧异常快照也可能某个源因为内部调度延迟行情到达时间滞后几百毫秒。如果你只有这一个源上层策略对错数据毫无感知一旦脏数据进入风控或者量化模型后果可能是资金级的事故。所以只要条件允许就不要只接一个源。多数据源的好处有多层。一是故障逃逸能力一个源断线了系统能立刻切换或自动加权到其他源。二是交叉验证能力可以通过多个源的比对来识别脏数据。三是盘口完整度补全不同源在某些合约上的深度、更新频率可能不同有的可能在夜间或特定时段不更新多源才能拼出完整的盘口画面。我自己见过最坑的情况是某第三源在个别合约上停止更新但数据源本身不报任何异常REST接口仍旧返回200这类故障不多接几个源很难发现。3.2 数据源分类与选型指标行情源按来源大体可以分成四类交易所直连、商业市场数据商、第三方聚合源、自采监测源。交易所直连的数据最权威、延迟最低但接入成本高、门槛严格通常是会员或者服务商级别才能拿商业数据商像Wind、Bloomberg这类接口丰富、字段规整但成本不菲且链路较长第三方聚合源一般是基于公开数据写的成本低、覆盖全但稳定性和准确度参差不齐自采监测源适合做辅助验证比如直接抓取公开页面的报价做低优先级参考。选型的时候可以从五个维度来打分延迟、完整性、稳定性、成本、兼容性。延迟决定实时性上限完整性看是否有缺量权限限制稳定性看历史宕机记录和维护团队响应速度成本包括接口授权、带宽、存储与人力成本兼容性则是字段、格式、时间标准和文化差异的适配难度。建议做一个简单的加权评分表比如延迟权重35%完整性25%稳定性25%成本10%兼容性5%决策前把候选源按真实数据和报价填进去别只看销售PPT。3.3 多源校准的算法落地多源校准的核心不在于代码多花哨而在于要清楚不同源的时间基准可能不一致。目前大家普遍采用的时间标准是交易所标准时间也就是每个数据源都会携带自己的源时间戳但跨源比较前不能直接用这个源时间戳对齐。正确做法是数据进入系统时先由行情处理器统一打上系统接收时间戳两套时间戳都保留。源时间戳用于业务判断接收时间戳用于链路延迟统计两者结合才能有效对齐多路数据。校准逻辑上我常用的方案是“多数源加权”加“异常源剔除”。假设有三个行情源A、B、C它们对同一个合约的最新价格分别给出三个值系统按每个源近期的延迟权重出自定义权值计算加权价作为当前基准价格。如果某一个源的值偏离加权均值超过预设阈值比如偏离在3%以上或者连续N帧值与中位数背离就自动标记该源状态为“SUSPICIOUS”暂时不参与加权等待连续若干帧恢复一致后再重新加入。需要注意阈值一定要按每个品种的动态波动率设定。豆粕和股指的波动率完全不同用固定阈值必然出乱子。在实际操作中我会给每个源建立“数据质量分数”。质量分数由延迟平均数、丢帧率、闪断次数、字段填充率几个指标汇总而来基准分100出现一次异常就扣分。当某个源的质量分数连续低于阈值系统自动发起重连、重启或者降级。这套机制比人工报警排查高效得多。3.4 源切换的灰度策略数据源切换比普通服务切换复杂因为行情源之间可能本身就存在几毫秒到几十毫秒的延迟差。全量直接切换客户端会瞬间看到价格跳动如果是算法交易平台是非常危险的。所以源切换一定要走灰度不能一刀切。我的做法是把客户端分组白名单组、小流量组、全量组。白名单组给自己团队和少量可信用户先用观察一段时间的数据差异小流量组放5%到10%的流量进来验证稳定性最后再放全量。切换期间新旧源输出并行进行对外分发层用同一份行情映射表业务端感知不到底层数据源变化。等新源运行超过约定的时间比如一个完整交易日质量分数稳定旧源才真正关闭。这个流程看着慢但它能避免在情绪压力巨大的交易日中做出不可逆的决定。4. 从0到1搭建行情网关一次完整落地复盘4.1 系统规模与硬件选型拿我自己做过的项目举例当时规模大概是15个数据源通道、每秒峰值60万tick、支持大约5000个在线客户端主要服务内部的策略团队和外部部分合作方。硬件上处理层用了三台16核64GB的物理机做双活部署配置基本相同机器之间通过内网万兆连接。接入层和分发层共用一组4核8GB虚拟机按通道和协议拆分成多个独立实例避免互相干扰。消息队列用Kafka搭了三个节点的集群每个broker分配独立的数据盘走SSD阵列避免磁盘IO成为瓶颈。这个规模下Java程序经过JVM和GC调优后可以稳定维持60万tick/s的吞吐分配4到6个线程做网络IO8个线程做业务处理剩下的线程池处理存储写放大。如果你们项目比这个规模大建议处理层直接切换成C或Rust等技术栈Java的内存模型在超高频场景下还是存在一定压力的。另外多说一句所有核心机器的内核网络参数比如socket buffer、backlog队列长度、文件描述符数量都需要提前调优不然连接数上来的瞬间就是雪崩。4.2 时钟同步这个隐形坑这个坑太值得单独拿出来讲了当时真是踩了整整两周才彻底搞清楚。多源校准过程中我一度发现明明所有数据源的接收时间戳都打在毫秒级了跨源比较时仍旧存在10毫秒以上的系统性偏差怎么调都调不平。后来才发现问题根本不在于校准算法而在于服务器之间的时钟不同步。NTP服务如果没配好不同进程计时用到的时钟源相差可能几十毫秒导致同一个物理时刻在不同机器上记录的时间戳差一大截这么一来任何精密的延迟统计都不可靠。解决办法就是在所有接入、处理、分发节点上统一配置NTP或chrony服务定期同步到同一个时间服务器源同时对NTP偏差设置告警。白天夜间的漂移要留意如果偏差超过50毫秒就要触发告警脚本自动纠正必要的话用PTP也就是精确时间协议做数据中心内部高精度同步。注意这个问题的排查很隐蔽因为它不会报错不会丢数据只会让整体统计和目标值差上一截太容易被人忽视了。4.3 常见问题排查速查表问题现象可能原因排查思路某通道延迟突然从2ms涨到500ms数据源本身故障或网络拥塞查看该通道最近消息的源时间戳与接收时间戳的差值定位瓶颈链路客户端连接全部闪断进程TCP backlog溢出或文件描述符耗尽检查ss -ln状态、ulimit配置、网络队列队列积压指标多源校准结果剧烈跳动某个源输出脏数据查看各源的质量分数标记掉异常源再观察校准价行情断流后恢复但数据有空洞主备切换后未做快照对齐检查是否先执行了快照同步再进入增量同步Kafka消费延迟飙升消费者处理能力不足或分区数不匹配查看消费组lag指标增加分区或横向扩容消费者实例收盘后有大量历史数据错位各源时间戳基准不一致统一用交易所源时间和接收时间双戳并检查NTP偏差这些问题是实时行情系统里最常遇到的几类。排查的时候记住一个原则先看时间再看序列号最后看业务逻辑。时间不对后面全白搭序列号不对数据就有洞等到业务层去追问题已经晚了。行情系统里排查时效性的问题跟调用链追踪不一样数据是流动的快速定位靠的是一层一层的指标检查而不是拍脑袋猜。最后分享一个我自己的小经验行情系统的监控不要只盯着延迟和丢包率这种技术指标更要关注“客户端侧的业务健康指标”比如策略侧的订单延迟、行情到策略进程的时间差、以及重启后恢复时间。某个底层指标异常不一定立刻影响上层但上层业务指标异常一定能在底层找到原因。当初我就是在一次凌晨的故障中靠着对比客户端业务指标发现是某个数据源提前进入了休市状态才意识到源状态检测不够灵敏后来补上了一个“行情源开闭市状态机”专门识别这种半死不活的源实测很管用。做实时系统保持敬畏心数据源和网络这些东西永远不可能百分百可控能做的就是尽量把每一个不可控环节都加上检测和兜底。