☰
Flink在证券行业实时数据分析中的实战经验与架构设计
2026/10/3 14:09:24 网站建设 项目流程

我最早接触Flink,其实是在证券系统的实时行情项目里。那时候团队折腾过Spark Streaming、折腾过自研的管道程序,最后才把Flink定为核心计算引擎。回过头来看,Flink在证券行业的落地,不是简单地换个计算框架,而是整个实时数据处理思路的升级——从“能算”到“算得准、算得起、算得快”。

这篇文章不打算讲教科书式的Flink原理,而是结合我在证券行业做实时市场数据分析的真实经验,聊聊业务怎么拆、架构怎么搭、代码怎么写、坑怎么踩。如果你正准备用Flink处理行情数据、交易流水、风险指标,或者正在纠结“到底该怎么把实时计算落地到自己的业务里”,这篇文章应该能给你一些参考。

1. 证券行业的实时数据分析,到底在分析什么

1.1 从业务场景倒推技术需求

证券行业的“实时市场数据分析”,听起来抽象,拆开来看其实就是几类非常具体的场景。

第一类是行情数据的实时加工。比如Level-1行情、Level-2行情,几十毫秒一条数据,包含最新价、成交量、买卖十档、逐笔成交。这些数据不是拿来看个数字,而是要经过各种计算,变成指标。比如实时涨跌幅、换手率、量比、资金流向,再比如技术指标MACD、KDJ、布林带,甚至机构自己定义的因子。行情数据本身是“原料”,Flink要做的是把原料实时加工成“半成品”甚至“成品”。

第二类是交易行为的实时监控。证券行业对风控的要求极高,自营、资管、做市等业务都需要实时盯盘。比如某个自营账户的持仓市值是不是突破了限额,某只股票的累计买入量是不是超过了监管阈值,某个策略的日内亏损是不是触及了止损线。这类监控要求毫秒级到秒级的响应,数据源往往是交易系统报回来的委托流水和成交回报。

第三类是实时统计分析。比如统计全市场或某个板块的实时成交额排名,统计北向资金的实时净买入,统计主力资金的日内流向。这些数据会被推送到大屏、APP、数据终端,供交易员、分析师、甚至普通投资者参考。

这些场景有个共同点:数据量大、时效要求高、计算逻辑复杂、结果不能错。行情数据算错一个数,风控指标差一分钱,带来的可能就是实实在在的损失和责任问题。

1.2 为什么偏偏是Flink

我之前用Spark Streaming做过类似的项目,最大的感受是“批里有流,流里有批”——微批的模式天然有延迟,一两秒的调度开销在实时行情场景里格外扎眼。而且Spark Streaming的exactly-once语义做起来很费劲,一旦涉及到和外部存储的交互,很容易出现重复或丢失。

而Flink天生就是流式计算引擎,它的事件驱动架构、精确一次语义、原生流处理能力,决定了它在延迟、准确性、状态管理上比Spark Streaming更有优势。特别是证券行业这种对准确性极其敏感的领域,Flink的checkpoint机制和端到端一致性保障,是它能够被信任的关键。

还有一个很现实的因素:Flink对时间语义的支持太完善了。事件时间、处理时间、摄入时间三种时间语义可以自由选择,配合watermark机制处理乱序数据,这对证券场景来说几乎是为量身定做的。行情数据在网络上传输,必然会有乱序、有延迟,能不能正确处理这些乱序数据,直接决定了计算出来的指标准不准。

所以团队的结论很明确:新项目直接用Flink,不再考虑其它方案。

2. 整体架构设计与组件选型

2.1 实时数据流的完整链路

证券行业实时数据分析的架构,说白了就是一条数据流水线:数据源 → 消息队列 → Flink计算 → 结果存储 → 应用展示。

数据源主要是两类:一类是交易所行情源,经过券商自己的行情网关解析后,生成统一的行情数据对象;另一类是交易系统产生的交易流水,通过日志或数据库binlog的方式对外发送。在项目里,这两类数据最终都进入了Kafka。

为什么中间要加一层Kafka而不是让Flink直接对接数据源?主要是为了削峰填谷和故障隔离。行情数据在开盘时段峰值极高,一秒几十万条是常有的事,如果Flink直接消费数据源,一旦Flink做checkpoint或者重启,数据源很容易被反压拖垮。有了Kafka做缓冲,Flink的消费速度可以自主控制,数据源侧只需要稳定地往Kafka里写就行。

Flink计算层承担了核心的加工逻辑,包括清洗、转换、指标计算、规则匹配、窗口聚合等。计算完的结果会有不同的去向:实时指标写入Redis供前端查询;明细数据写入ClickHouse做即席分析;告警事件写入ES并触发通知;汇总数据写入关系型数据库用于事后对账。

这里有个容易忽略的设计点:所有结果数据都必须带上数据时间戳和计算时间戳两个字段。数据时间戳是行情本身的发生时间,计算时间戳是Flink算完落库的时间。有了这两个字段,事后排查数据延迟、定位计算结果差异时会轻松很多。这个习惯我一直延续到现在。

2.2 资源规划与并行度设置

再说说Flink集群的部署。证券公司的IT环境一般比较敏感,很多系统要求在内网独立部署,不太可能直接用云上的托管Flink服务。所以我们的方案是在内网搭建独立Flink集群,用Flink on YARN的模式运行。

资源规划上,我踩过的教训是:不要一开始就追求大并行度。我们项目刚启动时,集群一共给了40个slot,我直接按最大并行度把作业跑了起来。结果每个slot上分配的TaskManager内存不够用,频繁Full GC,作业反复重启。后来改成按数据量倒推并行度:每个并行实例每秒处理5000条左右的数据是比较舒服的状态,按这个标准反推并行度,再预留20%的余量,跑起来就稳多了。

并行度设置还有一个原则:source、keyby、sink各段的并行度要分开设置。不要全局只设一个并行度,因为Kafka消费的并行度和下游ClickHouse写入的并行度往往不在一个量级。比如Kafka分区是12个,source并行度设为12;中间算子按key分布,可能需要24个并行度来避免数据倾斜;而ClickHouse写入是批量写入,6个并行度就够了。分开设置,资源利用率会高很多。

2.3 时间语义和数据准确性保障

在证券场景里,计算准不准是第一位的。Flink本身支持exactly-once语义,但真正要做到端到端精确一次,还需要上下游配合。

Kafka侧,我们开启了幂等生产者,同时把acks设置为all,确保消息不丢;Flink侧,开启checkpoint,interval设置为60秒(这个后面会细说为什么不是更短);Sink侧,写入ClickHouse时采用了去重表引擎,让数据库层面兜底防重。通过三层保障,基本做到了数据不重不丢。

时间语义的选择上,行情计算类的作业用事件时间,风控类的作业用处理时间。为什么这么分?行情指标必须按照业务时间的先后顺序来算,所以必须用事件时间;而风控规则讲究的是“现在立刻判断”,用处理时间才是最及时、最准确的。如果反过来,风控用事件时间,一旦某些事件晚到了,止损判断就会滞后,这是不能接受的。

3. 核心实现细节与实操代码

3.1 行情数据的接入与清洗

行情数据接入Flink,第一步是写Kafka消费者。这里有个很实用的技巧:不要直接用FlinkKafkaConsumer读原始byte[],而是让消息生产方在Kafka里存JSON,Flink侧统一用自定义的DeserializationSchema解析成POJO。

理由很简单:证券系统的数据字段非常多,一个行情对象动辄上百个字段,如果全程操作JSON字符串,每次计算都要做序列化和反序列化,性能损耗非常大。转成POJO之后,后续所有算子操作都是内存对象的属性访问,性能至少提升30%。

来看一段我们项目里的消费者示例:

DataStream<StockQuote> quoteStream = env.addSource( new FlinkKafkaConsumer<>( "stock_quote_topic", new StockQuoteSchema(), kafkaProps ) ).assignTimestampsAndWatermarks( new BoundedOutOfOrdernessTimestampExtractor<StockQuote>(Time.seconds(5)) { @Override public long extractTimestamp(StockQuote quote) { return quote.getTimestamp(); } } );

这里设置了5秒的乱序容忍度,意思是比当前最大事件时间晚不超过5秒的数据仍然会被纳入窗口计算。为什么是5秒而不是1秒?根据我们的线上数据统计,行情数据从交易所网关到Kafka的端到端延迟,99%都在3秒以内,5秒的冗余能在数据完整性和实时性之间取得平衡。

清洗逻辑也在这个阶段完成。比如过滤明显异常的行情数据——价格小于等于0、买卖档位价格倒挂、时间戳超过当前系统时间等,这些脏数据如果进入指标计算,会直接拉偏计算结果。清洗规则单独抽象成一个FilterFunction,方便后续增改。

3.2 实时指标计算的窗口设计

证券行情分析里最常见的一类需求是:计算过去N分钟某个股票的成交均价、涨跌幅、资金净流入等指标。这类需求在Flink里对应的就是滑动窗口。

不同的业务指标,窗口大小和滑动间隔完全不同。举个例子,资金流向指标通常按分钟级计算,需求是每分钟输出一次最近5分钟的累计结果,所以我们使用了滑动窗口,窗口长度5分钟,滑动间隔1分钟。

DataStream<FundFlow> fundFlowStream = quoteStream .keyBy(StockQuote::getStockCode) .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1))) .aggregate(new FundFlowAggregate()) .name("fund-flow-window");

.aggregate()比.apply()性能要好很多,因为它是增量计算,每个元素到达时会立即更新累加器状态,窗口触发时才输出最终结果。而.apply()需要把窗口内的所有元素缓存起来,窗口触发时再全量计算,内存开销和时间开销都更大。在行情这种高吞吐场景下,任何多余的计算都要避免。

这里还有一个细节:对于那些计算复杂度特别高的指标,比如盘中MACD这种需要依赖前值递归计算的指标,窗口聚合就不够用了,需要自定义一个有状态的ProcessFunction。用ValueState保存前一个周期的收盘价和EMA值,每个行情快照到达时增量更新,既保证了实时性,又不像窗口那样需要在内存里保留大量历史数据。

3.3 从MySQL实时同步到ClickHouse的实践

在整体需求的推进中,我们把一部分非核心但常用的数据从MySQL同步到了ClickHouse。比如客户持仓快照、历史交易记录、交易日历等,这些数据变更不是很频繁,但查询频率极高,放在MySQL里扛不住业务侧的并发查询。

这个需求在团队里由我负责,我用Flink CDC实现了全量加增量的同步方案:

CREATE TABLE mysql_orders ( order_id BIGINT PRIMARY KEY, stock_code STRING, client_id BIGINT, order_price DECIMAL(10, 2), order_qty INT, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '192.168.1.10', 'port' = '3306', 'username' = 'flink_user', 'password' = '******', 'database-name' = 'trading', 'table-name' = 'orders', 'scan.startup.mode' = 'initial' ); CREATE TABLE clickhouse_orders ( order_id BIGINT PRIMARY KEY, stock_code STRING, client_id BIGINT, order_price DECIMAL(10, 2), order_qty INT, update_time TIMESTAMP(3) ) WITH ( 'connector' = 'clickhouse', 'url' = 'jdbc:clickhouse://192.168.1.20:8123', 'table-name' = 'orders' ); INSERT INTO clickhouse_orders SELECT * FROM mysql_orders;

scan.startup.mode = 'initial'表示先执行全量快照,再自动切换到增量binlog同步。上线时这个配置帮了大忙——不需要手动预置数据,整个同步流程是自动完成的。

实际运维中遇到过一个最头疼的问题是:MySQL源表的字段类型和ClickHouse的字段类型如果不匹配,同步作业会静默失败。比如MySQL的DATETIME默认带毫秒,ClickHouse的DateTime精度只有秒,写入时数据直接被截断。后来我在Flink SQL里显式做了CAST转换,才彻底解决这个问题。

3.4 Spring Boot如何管理Flink作业

项目里还有一个公共模块:作业管理平台。我们用Spring Boot搭了一个轻量级的Web服务,负责Flink作业的提交、停止、重启、状态查询和日志查看,避免每次都要登录集群敲命令行。

和Flink的交互方式有两种:一种是通过Flink REST API直接管理作业和检查点;另一种是封装Flink SQL客户端,把SQL文本通过JDBC提交到Flink集群。我们选的是后者,因为团队里很多人对SQL更熟悉,用SQL描述数据处理逻辑比用Java写一套实现要直观得多。

Spring Boot整合Flink有个要注意的地方:不要在你的Spring Boot进程里启动Flink任务作为本地线程执行。这样做任务确实能跑,但是任务的容错性、资源隔离性、并行扩展能力会大打折扣。正确做法是Spring Boot作为客户端,通过flink run命令或者Flink REST API将任务提交到独立的Flink集群上执行。两者职责分离,Spring Boot只负责编排和展示,Flink集群负责真正的计算工作。

4. 常见问题与排查思路实录

4.1 Flink JDBC连接器抛异常的那些事儿

用Flink的JDBC连接器连接MySQL或者ClickHouse时,很多人会遇到Failed to deserialize parameter、Connection is not available这类的异常。我在项目里也遇到过,第一次处理时排查了很久。

先说Connection is not available。这个异常通常是连接池参数配置不当导致的。Flink的JDBC连接器默认连接池比较小,如果写入并发高,连接池会被耗尽,新请求就只能等待,等待超时了就报这个异常。解决办法是在WITH参数里调大连接池上限,同时设置合理的超时时间:

'sink.max-retries' = '3', 'jdbc.connection.max-retry-timeout' = '60s'

再说Failed to deserialize parameter。这个异常多半是类型映射问题。比如MySQL里的DECIMAL(20, 4)映射到Flink的DECIMAL精度不一致,ClickHouse的Nullable(Float64)映射到Flink的Double没问题,但反过来就可能出错。我的建议是:在Flink SQL里对JDBC数据源的所有字段都显式声明类型,不要依赖默认映射。

4.2 数据倾斜的经典场景:大单拆分与热点股票

实时市场分析里有一个天然的倾斜源——大单和热点股票。全市场几千只股票,但某一时刻可能只有几十只股票的交易量特别大,当按股票代码进行keyBy时,那几十个key的处理压力会明显高于其它key。

我们遇到过一个严重案例:某个作业并行度96,但CPU利用率最高的那个子任务已经打满了,其余95个只用了不到20%。数据倾斜直接导致整个作业的反压,从sink一直传递到source,Kafka消费 lag 急剧攀升。

针对这个场景,我采取了两种手段:

第一种是两阶段聚合。先按“股票代码+随机盐”做局部聚合,再按真正的股票代码做全局聚合。这个方法在统计维度不需要保留原始明细时非常有效。但注意,如果计算要求精确的TopN或去重,就不能用这个方案。

第二种是动态调整并行度。把热点key的数据(比如成交量超过阈值的股票代码)单独分流到一个高并行度的计算链路上,非热点key走另一个低并行度的链路。这个方案稍微复杂一些,但对热点集中型的场景效果最好。

4.3 Kafka消费Lag飙高后的排查流程

Flink作业的Lag飙升,在证券场景里往往不是Flink本身的问题,而是下游存储变慢了。我这边实际遇到的案例是:ClickHouse某个MergeTree分区因为合并操作异常积压了大量临时分区,查询和插入都变慢了,Flink的JDBC Sink写入阻塞,checkpoint超时,反压一路传导到Kafka消费者。

排查流程可以总结成一套固定的思路:先看Flink UI上的反压情况,确定哪个算子是瓶颈;再看瓶颈算子的输出指标,确认是计算慢还是写入慢;最后针对不同的原因采取不同的措施。这套排查流程在处理了十几个案例后,已经成了我们团队的标准应急预案。

4.4 Checkpoint超时的常见原因和应对策略

Checkpoint超时是我做Flink作业运维时遇到频率最高的告警之一。在证券行情场景下,checkpoint超时的常见原因有三种:

第一种是有反压。checkpoint barrier要在整个数据流中穿行,如果某个算子被反压堵住了,barrier传不到source,checkpoint就一直无法完成。这种情况要先解决反压。

第二种是状态太大。Flink在做checkpoint时需要把状态快照持久化到外部存储,状态太大,持久化时间就长,容易超时。我们上线初期用RocksDB作为状态后端,然后开启增量checkpoint,状态提交时间明显下降。

第三种是外部系统交互太慢。如果你的算子里有同步调用外部服务的逻辑,checkpoint时会有额外的对齐开销。优化方式是改成异步I/O(AsyncFunction)或者在状态里缓存数据、批量发送。

这个部分值得多说一句:不要把checkpoint interval设置得太短。不少人追求极致的故障恢复精确度,把interval设为10秒甚至5秒,结果checkpoint过于频繁,系统性能大幅下降。我实测下来,对于行情数据场景,60秒的interval既能保证恢复精度,又不会对性能造成明显影响。

5. 性能优化和稳定性保障的实战心得

5.1 状态后端选型:HashMap还是RocksDB

Flink的状态后端选择,直接影响作业的性能和稳定性。在做实时市场数据分析时,不同作业对状态的需求完全不同。

行情指标计算类作业,状态数据量通常不大,以窗口内部状态为主,用HashMap状态后端就够,读写速度极快,性能最好。但要注意,HashMap状态后端把所有状态都存在堆内存里,如果状态量大,GC压力会非常大,甚至OOM。

风控类作业,比如实时持仓监控,需要保存每个账户每只股票的累计交易量,状态量可能几百GB甚至更大,这时候必须用RocksDB状态后端。RocksDB将状态存储在本地磁盘,通过内存缓存提升读写性能,能够在有限内存下支持超大规模状态。

选型建议就一条:先估算状态量,再选状态后端。别凭感觉,也别图省事。

5.2 大状态作业的容灾与恢复

大状态作业的容灾是整个证券实时系统里我最看重的部分。

我们有一个持仓风控作业,状态里有全公司所有自营账户的持仓明细,加起来有近200GB的状态数据。这个作业如果被kill,恢复时间要花掉将近40分钟。在交易时段内这40分钟是致命的——风控规则全部失效,交易风险敞口完全暴露。

为了压缩恢复时间,我做了几件事:第一,开启RocksDB的增量checkpoint,每次checkpoint只上传变化的部分,checkpoint耗时从几分钟压缩到十几秒;第二,开启本地状态恢复(state.backend.local-recovery),让Flink在重启时优先从本地磁盘加载状态,避免每次都要从远程拉取全量;第三,把作业以session模式跑在独立资源队列里,避免多个大作业竞争恢复资源。

这套组合优化做完后,这个作业的最坏恢复时间从40分钟降到了3分钟以内。对于交易系统来说,3分钟的可接受程度比40分钟高太多了。

5.3 网络与内存参数的调整建议

证券公司的内网环境比较特殊,有时会开启各种安全策略,Flink在大流量下会遇到奇怪的网络问题。我遇到过TaskManager之间数据传输超时、反压检测误报、甚至在checkpoint时因为网络抖动导致barrier对齐失败。

针对这些场景,网络超时参数务必要根据内网的实际情况调整:

  • 如果数据量大,适当调大taskmanager.network.memory.min和taskmanager.network.memory.max,避免网络内存成为瓶颈。
  • 如果网络偶尔抖动,调大taskmanager.network.request-backoff.max,避免瞬时网络问题导致作业失败。
  • 如果频繁报Connection refused,检查TaskManager的端口范围是否被防火墙拦截。

内存参数方面,一个容易踩的坑是TaskManager的JVM Heap与Flink管理的堆外内存之间的关系。Flink的TaskManager内存分为框架内存、任务内存、网络内存和管理内存,如果你只设置taskmanager.memory.process.size而不调整各个子部分的配比,默认配置往往和你实际作业的需求并不匹配。比如RocksDB状态后端需要较多的管理内存,如果管理内存配小了,RocksDB会频繁刷盘,性能急剧下降。

我的建议是:用Flink的内存模型配置工具先估算一遍,再根据作业实际运行情况微调。不要直接拿社区的默认配置就上生产。

6. 给新上手的人一些实用建议

6.1 先做需求梳理,再写代码

我见过太多人一拿到需求就打开IDE写Flink代码,写到一半才发现核心逻辑被误解了。实时数据分析项目最花时间的地方不是写代码,而是把业务指标的计算口径理清楚。

比如“实时资金流入”这个指标,不同的人可能有不同的理解:是按主动买盘的成交量计算,还是按大单的净买入计算?是包含集合竞价还是只算连续竞价?是复权口径还是不复权?这些口径问题,如果不先和业务方确认清楚,做出来的结果一定是不被认可的。

我在项目启动时专门花了两周时间,和业务团队逐条梳理了所有指标的计算口径,写成了一份需求文档。也正是这份文档,成为了后来验证计算结果的依据。这份前置工作,远比多写几行代码重要得多。

6.2 从简单场景开始练手

Flink的学习曲线比较陡,如果你是完全的新手,我建议不要直接挑战复杂的多阶段计算。从最简单的“读取Kafka → 过滤 → 写入Redis”开始,先把环境跑通,把Flink作业的生命周期搞清楚,再去尝试窗口、状态、checkpoint这些复杂特性。

练手时最容易遇到的坑反而是环境搭建。Flink本身是一个分布式系统,需要配置JobManager和TaskManager,需要处理HDFS或者S3等外部存储的依赖。这部分环境问题会消耗大量时间,建议直接把本地IDE调试和Flink SQL客户端结合起来,用SQL的方式快速验证逻辑,再用DataStream API做细致的表演级开发。

6.3 监控和告警:把作业当成系统对待

Flink作业上线了,不代表就结束了。在证券行业,昨天还在正常跑的作业,今天可能因为数据量暴增而反压,可能因为上游Kafka topic被误删而持续重试。所以监控和告警必须从第一天就开始建设。

至少要监控这些指标:作业是否Running、checkpoint是否成功、Kafka消费Lag、每秒处理条数、反压比例、状态大小变化。这些指标全部接入告警,达到阈值就立刻通知到人。

我之前带的团队里流传一句话:没有监控的实时系统,是定时炸弹。你永远不知道它会在什么时候出问题,但只要出了问题,损失就已经造成了。

6.4 设计数据回补机制,给自己留退路

行情数据是高价值、不可再生的数据,一旦因为系统故障丢掉一段,补都补不回来。所以我在设计实时链路时,永远不会只有一条流,而是会同时保留一个“原始数据落盘”的旁路。Kafka里的原始行情数据同时被写入HDFS或对象存储,作为离线备份。

这样做的原因是:当实时计算因为各种原因出现bug,计算结果已经错了的时候,我们可以用离线存储的原始数据重新计算一遍,然后用计算结果修正实时链路的数据。这在事后对账和指标修复中非常有用。

这个机制在证券行业特别重要,因为监管要求交易数据可追溯、可审计,数据出问题必须能够回补。实测下来,多一份原始数据的备份,并不会占用多少存储成本,但万一出了事,它能救命的。

复盘一下我个人最深的体会

做了几年证券行业的实时数据分析和Flink开发,我的一个深刻感受是:技术本身从来不是最大的挑战,搞清楚业务需求和组织协同才是。

Flink的能力已经足够强大,实时计算在证券行业的应用也已经很成熟,但真正让项目顺利运转下去的,往往是这些看不见的功夫:前期和业务团队反复对齐的计算口径,中期为容灾恢复做的那些“多余”设计,后期在问题排查中沉淀下来的一套标准流程。

如果这篇文章能给你带来一个启发,我希望是:在动手写Flink代码之前,先花足够的时间把业务想清楚;在作业稳定运行之后,也不放松对监控和容灾的要求。实时系统的上线只是起点,稳定可靠地运行才是真正的考验。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询