ARTICLE DETAIL

资讯详情

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

证券海量数据秒级分发:Flink实时计算链路全解析

证券海量数据秒级分发:Flink实时计算链路全解析 我在券商和行情服务商两头都待过做实时计算和海量数据分发这块快十年了。今天想复盘一个让我印象很深的项目证券交易所海量金融数据秒级分发的实时计算落地。这类项目外面讲得玄乎其实拆开看就是一条数据分发链路能不能在极短的时间里把海量行情、订单、成交数据稳定送到每个需要的人手上。难的不是某个组件不会用而是整条链路在每天几个特定时间窗口内要同时扛住每秒几十万条增量、保证顺序、保证不丢、保证秒级触达所有订阅方还要在出故障时能自证清白。这个项目的背景不算特殊上游对接交易核心的实时数据流下游有行情大屏、柜台系统、量化策略引擎、风控模块等几十个订阅方。不同订阅方对数据的要求还互相打架——有的要全量逐笔有的只要自己关注的几十只标的有的能容忍几百毫秒延迟有的恨不得拿到毫秒级推送。整个项目做完我最深的感觉是实时计算落地不是写几个Flink作业就完事真正的工程量在链路设计、故障排查和一致性保障上。这篇文章我尽量把能复现的思考过程写出来包括早期方案怎么踩坑、实时计算层每个关键决策的理由、以及上线后最难啃的几个故障排查过程。1. 先搞清楚证券交易场景的海量数据到底长什么样1.1 三类核心数据流的时效性差异证券交易的实时数据流并不是单一一种至少可以分成三类性质完全不同。第一类是行情数据包括快照行情和逐笔成交。量最大、峰值最高一天下来几百亿条非常正常。行情数据允许偶尔的乱序但不能长期断流断流意味着行情大屏直接卡住量化策略拿到的是残缺数据。第二类是委托和撤单回报。这类数据对顺序极其敏感。一个委托的状态机通常是已报→部分成交→全部成交或者已报→已撤如果撤单回报先到、委托回报后到下游柜台系统就可能把资金状态算错明明冻结的资金被提前释放。这类数据单量不如行情大但它的正确性要求比性能要求高得多。第三类是成交回报和风险指标。成交回报直接关系到清算和对账必须精确、可回溯。而风险指标是实时计算层加工出来的结果比如某个账户的当日累计买入金额、持仓集中度、撤单率这类指标往往要在一分钟内滚动更新给风控系统做实时拦截。这三类数据混在一起的时候不要想着用一个统一通道搞定。最典型的设计是把它们拆成独立的逻辑链路共享底层基础设施但隔离故障。我们当时做的时候把行情类数据和高频回报类数据分别接入不同的Topic和计算拓扑原因很简单两类数据的流量模型、消费方式、失败后果完全不一样。混在一起一次行情洪峰就可能把委托回报的通道冲垮这在真实事故里是发生过很多次的。1.2 秒级的全链路预算秒级分发这个说法听起来不紧不慢但在真实的金融系统里它是被拆成一段一段毫秒预算去抠的。我当时给团队画过一张时间分配表把从数据产生到消费端收到数据的全链路切成五段数据产生与出口编码行情网关侧编码和发送预留不超过30毫秒采集层接收与写入消息队列预留50到80毫秒实时计算层处理包括清洗、规整、指标计算和路由预留100到200毫秒分发通道传输预留20到50毫秒消费端接收解码预留50到100毫秒。整体加起来正常情况下P50应该在300到500毫秒之间P99要控制在1秒以内这样对外才能承诺秒级。很多人以为秒级就是最终1秒内到就行这是最大的误区。端到端1秒意味着里面有大量环节不能超过几十毫秒。比如架构里一旦出现某个环节偶尔抖动到300毫秒那么P99就很难看。所以做这类项目第一件事不是上线而是先把每段链路的延迟标准定清楚监控和报警按这个标准去设定。否则你根本说不清延迟到底卡在哪一段。1.3 海量的真实体感再说说海量这个词。没真正处理过的人可能觉得每秒几万条已经很夸张了但场内交易场景里行情峰值可以到每秒几十万条消息而且不是均匀分布是集中在几个极端的时间点上。我们实测过开市前集合竞价结束到连续竞价刚开始的那一分钟逐笔消息量是平时均值的几十倍。这种爆发式的流量模型直接决定了架构设计里必须有削峰、缓冲、降级的手段。用纯粹的同步推送任何一层都吃不住用异步队列又得处理延迟和积压的代价。所以最终架构都是异步削峰批量合并退避淘汰的组合后面我会详细讲我们是怎么平衡的。2. 早期方案复盘Kafka直推为什么在峰值冲击下撑不住2.1 最初的链路形态项目刚启动的时候我们用的是很多团队第一反应会选的架构采集网关把数据写入多个Kafka Topic各个消费方按需订阅自己拉数据、自己处理。Kafka吞吐高、生态好、可回溯看起来完美契合。实际跑下来问题非常快暴露。最表象的问题是每个消费方都得自己写一套解析、清洗、补全逻辑。行情大屏的团队写了一遍字段映射量化策略团队又写了一遍风控团队再写一遍。三套代码对同一个字段的理解还不一致比如成交金额到底是含手续费还是不含A股和B股的精度保留几位每个团队都有自己的口径。于是同一份数据源下游算出来的指标经常对不上出了数据争议根本没法定位是谁的问题。更深层的问题是Kafka直推本质上没有计算层只有传输层各干各的消费层。这种架构里消费者被直接暴露在数据洪峰前面没有任何东西帮它做过滤、聚合、路由。一个订阅了全市场行情的客户端在峰值时要面对每秒几十万条消息它内部的序列化、推送、落库逻辑只要有一个环节处理不过来lag就会迅速堆积堆积了又传染Kafka broker端的磁盘和网络。2.2 分区热点与数据放大效应Kafka直推还有一个在证券数据里特别要命的问题分区热点。我们知道Kafka的并行度依赖分区但行情数据天然是热门标的高度集中。一只活跃股的逐笔成交可能是冷门股的几百倍按标的做Keyed分区的话那只股票所在分区的数据量会远远超过其他分区。最开始我们用股票代码做Key结果某个热门标的把单个分区的消费者线程打满其他分区闲得不行。这时候你会看到Kafka消费组里的Lag极度不均单个消费者线程CPU 100%另外几个只有百分之几。出现了热点分区整个链路吞吐就被单分区上限卡住这是典型的木桶效应。还有一个放大效应行情数据如果每个消费方各拉一份同一份数据在Kafka里被反复消费broker出口带宽和下游网络开销成倍增长。我们算过一次下游有几十个消费组每一条逐笔消息至少被复制几十份加上跨机房间的流量成本极其夸张。2.3 实时计算层的切入点Kafka直推的这些问题最后逼着我们引入实时计算层。它的定位不是取代Kafka而是夹在Kafka和最终消费端之间负责四件事统一清洗和标准化把上游各种格式的报文变成一套内部数据模型下游所有人拿到的字段、精度、类型完全一致按数据分类做聚合和过滤行情类数据可以做合并委托撤销类数据原样透传风险类指标在窗口内计算完再下发智能路由和分组分发根据订阅关系把全量数据拆成不同订阅集每个订阅方只拿到自己需要的部分幂等标记和序列号生成在上游报文基础上打上内部全局序列号为下游对账和去重提供依据。实时计算层用什么引擎当时我们对比过Flink、Spark Streaming和自研的流处理框架。Spark Streaming的微批模式在秒级延迟上天然吃亏自研框架维护成本太高最终选了Flink。原因很实际它的低延迟流处理能力、成熟的状态管理、Exactly-Once语义、以及丰富的连接器能让我们把精力集中在业务逻辑而不是底层传输上。后面所有内容都基于Flink这个选型展开。3. Flink实时计算层的工程化细节算子、状态、背压逐个说3.1 窗口和水位线的取舍实时计算层里第一个要过问的问题是哪些计算需要窗口窗口怎么开。证券数据的典型窗口计算包括1分钟级的价格统计、累计成交量、涨速、换手率等指标。我们的设计里大部分指标窗口是基于处理时间做的滚动窗口因为业务上盯的是现在这一刻的状态而不是按事件时间严格对齐。但这里有个细节必须处理行情数据从采集到进入Flink可能有十几毫秒到几十毫秒的延迟个别情况下网络重传会把它拉到几百毫秒。如果完全按处理时间算迟到的消息会被切到下一个窗口导致指标跳变。所以我们最终采用的是处理时间为主、事件时间校验的方案核心指标用滑动窗口加低延迟触发同时把上游报文的原始时间戳记录下来对账时用事件时间做校验发现偏差超过阈值就报警。这种双轨方案比单纯纠结用哪种时间语义要实用得多。纯事件时间在证券行情这种超高吞吐场景下会引入Watermark和数据积压的麻烦纯处理时间又无法还原事实顺序。双轨的好处是实盘快、可追溯代价是实现复杂度多一些但对金融项目是值得的。3.2 状态存储与去重关联的具体做法Flink的状态管理是我们项目里另一个重点。当时我们有三类典型的带状态计算第一是去重。上游采集网关在某些异常恢复场景下会重发报文同一个委托号可能出现多次。我们用KeyedState里面的ValueState或者MapState以业务主键比如委托号、成交编号记录最近一次处理过的消息ID和时间戳配置TTL超过一定时间就自动清理。这里要特别注意TTL的粒度太短了达不到去重效果太长了状态无限膨胀。我们根据业务场景定了不同的TTL行情类消息保留30秒委托回报类保留24小时成交回报类保留7天。第二是关联补全。行情快照和逐笔成交经常需要关联证券基础信息比如股票代码对应的名称、板块、涨跌停价格。我们把这些基础信息做成广播状态BroadcastState每日开市前加载一次盘中如果有零星变更走广播流更新。用广播状态的好处是避免了每个计算任务都去查外部存储大大降低延迟和外部依赖压力。第三是流式聚合指标。风控类的实时指标比如某账户最近1分钟撤单率需要在一个滚动窗口内维护不断更新的状态。这里我们使用了RocksDB StateBackend原因是集群内存有限且状态规模可能达到几百GB。堆内存StateBackend虽然访问快但会导致JVM GC极其严重甚至Full GC卡顿几秒——这在秒级分发里是不可接受的。RocksDB虽然单条访问比堆内存慢但配合Flink的异步检查点和本地缓存整体表现稳定很多。3.3 背压、并行度和吞吐调优实时计算层落地过程中背压是绕不开的话题。Flink的背压机制会自动从下游反压到上游但如果不去干预它会导致整个拓扑的吞吐下降、事件延迟上升。我们在最开始忽略了一个重要问题Flink作业并行度设置成全局统一值但不同算子的负载完全不同。第一个优化手段是拆开算子链。默认情况下Flink会把能链在一起的算子合并成一条链减少序列化开销。但对于我们这种需要独立调优的作业不能无脑合并。行情解析算子的CPU消耗很高而路由下发的算子瓶颈在网络IO上把它们拆开单独设置并行度效果立竿见影。第二个优化是合理设置KeyBy的分区策略。前面踩过热点分区的坑在Flink里我们通过加盐salted key的方式解决对热门标的的Key后面拼一个随机后缀让它分布到多个子任务上处理完再按原始Key聚合。这个方法会带来一些网络shuffle的额外开销但峰值提升的收益远超这个代价。第三个优化是调整缓冲区和内存参数。Flink任务管理器的网络缓冲、RocksDB的块缓存、堆内存占比这些参数在默认配置下只能用在小作业上。我们结合上游每秒消息条数、每条消息平均大小估算出每秒需要处理的字节数再反推需要的TaskManager数量和内存配置。这个过程不是一次搞定而是通过压测不断调整。我们的经验是每台物理机上的TaskManager不要贪多8GB到16GB堆内存比较稳留足系统页缓存和RocksDB的堆外内存。4. 秒级分发链路的五大痛点与完整排查链路4.1 开市洪峰击穿分发通道第一个严重故障发生在实盘第二周。早上9点30分连续竞价刚开始监控大屏上每个消费端的Lag曲线像坐了火箭一样往上蹿从0跳到了几万随后就是雪崩行情大屏开始卡死策略引擎收到数据的时间戳落后了十几秒。我们的排查链路是这样的先看Flink作业本身指标显示计算层吞吐其实没有明显下降CPU和内存都在正常范围。接着看Kafka的消费Lag发现消费组Lag暴涨但Flink侧的消费速率并没有掉。继续往下追问题出在Kafka broker层。开盘瞬间的流量洪峰让broker的磁盘IO被打满写盘延迟从5毫秒飙到200毫秒Flink的consumer拉取请求全部排队消费速率被明显压制。根因找到了解法分三层第一层在采集侧加一个预聚合环节在源头把同一毫秒内同一标的的多条逐笔消息合并成一条减少下游需要处理的消息条数第二层是给Kafka配置限流和优先级把委托和成交回报这类高优先级消息放到独立的小Topic并配更高吞吐的Broker降低行情洪峰对它们的冲击第三层是在Flink和Kafka之间增加弹性缓冲允许一定程度的积压。最终开市洪峰不再直接打穿链路P99恢复到几百毫秒内。4.2 跨机房延迟长尾让秒级SLA失效上线几个月后我们开始做跨机房容灾把数据分发链路扩展到生产机房和灾备机房两侧。问题随即出现端到端延迟的P50很漂亮只有两三百毫秒但P99经常跳到2秒以上。这种长尾问题最折磨人因为大部分用户觉得还行但少数敏感客户已经在投诉了。排查过程很绕。我们先怀疑Flink的GC但GC监控很正常又怀疑Kafka的page cache命中率也正常。最后是通过抓网络包定位到是机房专线的TCP重传。行情数据量太大专线出口带宽被占满TCP拥塞导致随机重传部分消息被滞留在网络缓冲区里端到端延迟自然爆表。解决思路是双管齐下一方面把跨机房复制改成异步批量化压缩传输减少专线带宽占用另一方面在分发层设计就近优先策略每个消费端优先连接所在机房的本地通道只有本地通道故障才切到对端。这样既保证了容灾能力又不让日常流量都挤在昂贵的专线上。改完之后P99稳定在600毫秒以内。4.3 消费端能力不均衡拖累整条分发第三类问题出在分发通道的设计上。我们的分发层简单说就是一个订阅推送服务实时计算层把加工后的数据推给这个服务再由它分发给各个消费端。某天一个做行情大屏的订阅方突然反映大屏刷新延迟了接近10秒。刚开始以为是行情源问题但查了之后发现其他消费端都很正常。继续排查才发现这个订阅方一个客户端进程里开了上千个WebSocket连接每个连接都订阅了好几百个标的分发服务需要给每个连接逐条推送。更糟的是这些推送逻辑在同一个JVM线程池里执行其中一个连接的网络慢把线程池的线程全部占住其他连接的推送全部排队。这暴露了分发层的一个设计缺陷没有做订阅隔离和配额控制。我们随后重构分发服务核心改动是每个订阅方按连接池隔离线程慢消费者慢慢处理不能拖累别人支持多标的批量消息合成一个WebSocket帧里塞多条数据减少连接数和推送次数连接层面做背压消费端处理不过来时主动降速而不是无限制地累积消息。改完后同一个订阅方即使一个连接卡住也不会影响其他连接的秒级推送。4.4 乱序与重复消费金融场景最隐蔽的问题前面三类问题都是性能层面的第四类问题直接上升到正确性层面。某次灾备演练之后柜台系统反馈成交回报出现了重复个别撤单回报的状态顺序还错乱了一下。排查这个问题的完整链路非常费劲。一开始我们误以为是Flink状态恢复导致的重复因为故障恢复会触发Source重放。检查了Flink的checkpoint配置和外部系统Sink幂等性发现我们确实解决了计算层自身的重复也就是精确一次语义但问题出在更上游和更下游。上游的采集网关在灾备切换时会从本地日志里重放未确认的报文这本身就产生了重复。Flink即使做了一次消费去重也只能基于业务主键去重但如果主键在重放时被网关重新生成了去重就失效了。这要求我们把去重逻辑再往前推在网关出口就对报文打上唯一消息ID这个ID要跨故障切换保持不变。下游的乱序问题则是因为Kafka在灾备场景下涉及跨分区复制切换后新消费者看到的消息顺序可能和原来不一致。我们最后的解法是在分发层的入口为每一类消息维护一个全局单调递增的序列号消费方收到的数据必须按这个序列号校验发现跳号和乱序就触发主动补拉或告警。这相当于把顺序保证的责任从传输层上移到协议层成本更高但正确性真正可控。4.5 故障切换恢复时间不可控最后一个痛点是故障恢复时间。前面说了我们用RocksDB做状态存储跑一段时间后状态体量涨到了几十GB。有一次模拟主集群故障Flink作业开始迁移恢复结果大家等了一个多小时checkpoint都还没完全恢复。这在金融场景里是没有办法接受的因为一个交易时段的分发中断几分钟可能都会酿成事故。排查结论指向两个原因一是Checkpoint持续时间过长RocksDB在快照时需要复制大量SST文件二是恢复时要从检查点完整加载状态冷启动读盘消耗巨大。我们做了几个针对性的调整开启Flink的增量Checkpoint每次只上传变化部分检查点时间从几分钟降到了十几秒调整RocksDB的WriteBuffer和并行压缩参数减少快照期间的IO阻塞将大作业拆成多个子作业每个子作业状态独立失败时只恢复受影响的部分最后是增加StandbyTaskManager平时预热就绪发生故障时能立刻接受任务。这一整套组合拳打完我们的故障恢复时间从一个多小时压缩到分钟级总算达到了业务方的最低容忍线。5. 一致性、对账与灰度金融项目真正难的部分在性能之外5.1 端到端序列号机制性能对金融项目来说只是第一关过了性能关接下来是更磨人的一致性和可验证性。我们在完成分发链路的性能优化后才正式把序列号机制提升为全链路的强制要求。具体做法是在实时计算层的出口为每一条待分发的消息生成一个全局流水号前缀标识数据类型行情、委托、成交中间是业务日期后面是单调递增的数字。分发层在推送时携带这个序列号消费端在本地做单调性校验。正常情况下序列号连续递增如果发现跳号意味着中间丢了数据如果发现重复说明出现了重放。这两种情况都触发对账流程。这个机制解决了前面所有说不清、查不明的问题。它的设计关键是序列号的生成必须是单一分配器否则多节点并发生成会出现乱序和重复。我们当时为了性能尝试过分段发号比如每个节点拿一个号段但后来发现故障切换时号段管理很容易出洞最终还是改成了无状态的分布式发号服务通过批量预取方式把它对吞吐的影响降到最低。5.2 对账与补数通道序列号机制配上对账系统才真正让数据分发正确变成可证明的。我们对账分成三个层级分钟级抽样核对、日终全量核对、异常事件触发核对。分钟级抽样核对的做法是每个消费端每过一分钟上报本分钟内收到的最大序列号和消息条数快照对账服务比对分发层记录和消费端上报值如果最大序列号不一致说明存在丢失或错序。日终全量核对是在收盘后把全量日志重放一遍比对两端是否完全一致这能发现抽样覆盖不到的小概率问题。异常事件触发核对则是在故障切换、重启、网络重连之后立即自动启动一段专项校验确保恢复过程中没有任何差错。对账发现差异之后必须有补数通道。我们建设了一个历史回放服务可以根据消费端当前缺失的序列号区间自动把对应消息从Kafka日志中重新读取并补发。这个通道在设计之初被很多人认为是多余的投入但后来实际使用频率远超预期几乎所有消费端都依赖它解决过偶发漏数据的问题。5.3 双跑与灰度切量新链路上线前我们还做了和新老链路的双跑验证。那段时间实时计算层并行跑了两套一套是旧的Kafka直推链路一套是新的Flink统一分发平台链路。两边同时消费同一份上游数据把结果比对一致性。这个阶段大概持续了两个交易周中间发现了很多隐藏问题包括字段精度差异、时区转换、边界条件处理等。双跑通过后我们采用灰度切量的方式逐步放量。第一批放给内部监控和大屏第二批放给非交易核心的策略模拟环境第三批才放给真实柜台和风控模块。每一批都观察至少两个完整交易日的指标包括延迟分位数、消息命中率、对账差异数确认没问题才进入下一批。灰度切量最大的好处是如果新链路有问题影响面是可控的不会一次性把所有业务都拖下水。5.4 监控体系按SLO拆解最后必须说的是监控体系。实时计算做久了你会发现没有好的监控你根本不知道系统当前处于什么状态。我们的监控不是笼统的延迟、吞吐、成功率三大件而是按照SLO拆解成端到端可观测的指标链。每一段链路都单独埋点采集段的接收速率和最新延迟、Kafka段的写入延迟和消费Lag、Flink段的处理速率和事件时间延迟、分发段的推送队列深度和单连接推送延迟、消费端的ACK延迟和本地处理耗时。每个指标都设定阈值按下对称级别定义报警。最关键的是一个新鲜度指标比较数据中心最近一条消息的时间戳和当前时间如果这个差值持续超过几秒说明链路某处已经卡住。这比单纯看Lag更直观因为它直接反映业务视角的数据新鲜度。做了这套监控之后我们排查故障的时间从几小时缩短到十几分钟。很多时候告警一响直接定位到某一小段链路的指标异常方便太多。6. 写在最后的实战心得项目做到后期我越来越觉得实时计算在金融场景里的真正价值不是更快而是更可靠地快。性能指标再漂亮如果无法证明数据没丢、没重、有序金融客户是不敢用的。所以做这类系统从一开始就要把可观测性、对账能力、故障逃生通道当成一等公民来设计而不是事后补救。还有一点个人体会是实时计算不是只要搞懂Flink就够的。它是一整条链路的事从上游报文格式、采集网关的可靠性到中间计算引擎的优化再到分发层的流量隔离、消费端的幂等设计每一环都决定最终体验。整个项目里我们踩过的坑绝大多数不是Flink本身的问题而是链路中各环节的衔接问题。所以我建议正在做类似项目的朋友多花时间理清楚数据契约把全链路的时间预算和异常场景逐项列出再动手写代码后面会省非常多事。最后分享一个小技巧所有关键环节的日志都要带上全局消息ID和序列号。日常看着觉得日志量大、写得啰嗦但真正遇到客户投诉我这条数据为什么延迟了我这条数据怎么重复了的时候这些跨系统串联的ID就是唯一能救你的线索。这套东西值得在全链路建设的第一天就做进去。
返回列表