做证券行业的实时数据分析,最折磨人的不是写不出 SQL,而是数据一多一快,整个链路就开始原地翻车。开盘前一切正常,盘中行情像瀑布一样灌进来,迟到的消息、突发的尖峰、重复推送的 tick,全挤在同一秒到达。我在这条路上折腾了很久,最后稳定下来的核心引擎就是 Flink。Flink 在证券行业的实时市场数据分析里承担的是最关键的计算角色:从行情接入、乱序矫正、窗口聚合,到异动预警和风控指标计算,它把“实时”两个字真正落到了代码层面。这篇文章不打算写成官方文档,我会从架构设计讲到实战代码,再聊到生产环境里那些不得不防的坑,适合正在用或准备用 Flink 做实时行情的工程师朋友。
1. 先想清楚架构:实时市场数据分析的整体设计
1.1 为什么是 Flink:流式处理与证券场景的匹配点
证券行情本质上是一连串持续产生的事件流,不是一张静态表。价格、成交量、委托、成交,每一秒都在变。早年间很多团队用 Spark Streaming,它本质上是“微批”处理,把流切成一小段一小段去跑,秒级延迟可以做,但延迟有天花板;一旦碰到盘中行情瞬时暴涨,批间隔和调度开销都会放大延迟。我在选型时把延迟、乱序处理、状态管理、端到端一致性这四个点单独拎出来对比,Flink 在这四项上几乎都占全了:真正的逐条事件处理、事件时间配合 watermark 处理乱序、状态编程让跨事件计算变得自然、checkpoint 加两阶段提交做到精确一次。
打个比方,批处理像读完当天报纸再写摘要,流处理是直播间里的实时弹幕,每一条消息都值得立刻响应。Storm 也能做实时,但状态支持和时段窗口要自己造轮子,维护成本高;Kafka Streams 是轻量,但复杂的状态型风控和窗口分析在 Flink 里表达更成熟。证券行业里的实时性要求往往是百毫秒到秒级,Flink 不是唯一能做的,却是综合成本最低的一个。
当然,如果业务需求只是“每 5 分钟刷新一次汇总”,那用定时任务或者 Spark 也没问题,没必要强行上 Flink。选型要匹配需求,不能为了追框架。这一点在证券场景里特别重要,因为实时计算集群的资源开销并不便宜,架构越复杂,运维成本越高。
1.2 从行情源到分析结果:一条完整的实时数据链路
证券行业实时数据链路,比较通用的分层是:数据源 -> 消息队列 -> 实时计算 -> 存储与服务。数据源包括交易所行情、内部订单系统、MySQL 业务库;中间层基本是 Kafka,因为行情高峰可以达到每秒几十万条,需要它削峰、解耦、提供事后重放;Flink 跑在 Kafka 后面,负责所有需要“算”的逻辑;最后数据落到 ClickHouse 供分析师查询,Redis 放实时快照,MySQL 存结果和元数据。
为什么中间必须放 Kafka?我在行情高峰看过很多次,上游一抖动,后面没有缓冲层,Flink 消费并发一增大,直接把 Source 和 Sink 打满,整个集群雪崩。Kafka 相当于一个蓄水池:水位可以涨,只要不到坝顶就不会漫出来。做实时数据分析的人经常为了追求低延迟而把 Kafka 省略,这是最容易在证券场景翻车的地方。此外 Kafka 还能做多消费者组隔离,让实时计算和日志采集互不干扰,消息留存能力也让计算程序重启后可以安全回放。
存储层怎么选,我用一张表来总结:
| 组件 | 特点 | 适合存储的数据 | 选型理由 |
|---|---|---|---|
| ClickHouse | 列式存储,聚合查询极快 | 行情明细、下单明细、分钟级聚合结果 | 分析师即席查询,大数据量也能秒级返回 |
| Redis | 内存 KV,低延迟 | 最新价、涨跌停状态、热门股票排行 | 快照类数据,毫秒级读写 |
| MySQL | 关系型,事务支持好 | 基础资料、用户配置、最终结果报表 | 低频场景和事务一致性要求高的地方 |
| Elasticsearch | 倒排索引 | 日志检索、异常文本 | 有搜索或日志需求再引入,别一开始就全家桶 |
实际项目中还有一条常见支线:把 MySQL 的业务数据实时同步到 ClickHouse,供分析侧使用。比如用户信息、委托记录原本在 MySQL,但行情分析要跟它们做关联。Flink 可以同时承担这个同步工作,比传统的定时接入工具更快更灵活。后面我会专门讲这段用 Flink JDBC 连接器怎么实现。
1.3 四大核心场景拆解
证券行业里的实时市场数据分析,我把它拆成四类典型场景,每一类的技术重心差别很大。
实时行情快照:这是最基础的需求,把 tick 流处理完更新到 Redis,供行情 APP 或微服务查询。难点不在计算,而在时序控制。同一只股票的消息是严格递增的,但网络传输中可能乱序,一个旧的事件晚到了,不能覆盖掉新的快照。处理方式是用事件时间戳和序列号做去重判断,只更新更新的数据。
分钟级聚合:这是窗口用得最频繁的场景。每只股票每分钟的开盘、收盘、最高、最低、成交量、成交额,用滚动窗口就能实现。真正要做好的地方是增量聚合,不能在窗口内缓存全量明细然后再遍历,数据一大就 OOM。
异动预警:比如某只股票一分钟内成交量放大 20 倍,或者短时间价格剧烈波动。这类场景通常用滑动窗口,配合规则判断。规则简单时 ProcessFunction 直接写,规则复杂时用 CEP。赌用一个上游结果时,要额外小心重复计算和预警风暴。
实时风险监控:比如对单个账户过去 10 分钟内的委托频率计数,超过阈值进入告警状态。这类计算天然需要状态,Flink 把状态保存在本地,比每次去查询 Redis 快一个数量级。状态变大后放到 RocksDB,避免全部压在堆内存里把 GC 打爆。
2. 核心概念与环境准备:动手之前先扫盲
这一节既是给菜鸟补基础,也让老手回头查缺补漏。证券场景的特殊性决定了我们不能只看教程,还要理解数据特征对框架提出的要求。
2.1 流处理的时间、窗口与状态:证券数据计算的地基
很多人在初学 Flink 时被 watermark 劝退,其实换个角度理解就很顺。数据产生的时间叫事件时间,Flink 机器处理这条数据时的本地时间叫处理时间。证券行情里,我们关心的永远是事件时间,因为交易所时间戳才是业务事实。网络传输会导致数据晚到,假设某只股票上午 10:00:00 产生的一条 tick,可能到 10:00:06 才进入 Flink。如果按处理时间切窗口,这条数据就会落进错误的统计窗口里。watermark 的含义是“我保证事件时间早于此刻的数据已经全部到达”,Flink 依据它决定窗口什么时候触发计算。
窗口有三种基本形态。滚动窗口像设定好频率的闹钟,固定 10 秒一次;滑动窗口像每隔 5 秒回看过去 10 秒,适合“最近 5 分钟涨跌幅”这类重叠统计;会话窗口适合检测一段连续活跃行为,证券里偶尔用来识别连续盯盘和交易。大多数行情聚合用滚动和滑动就够,会话情况比较特殊。
状态可以理解成 Flink 程序在内存或 RocksDB 里给每个 key 维护一份“小账本”。证券里做风控计数时,如果用外部 Redis 要等一次网络往返,用状态就天然本地化,快一个数量级,代码写起来也更加直接。关键是搞清楚 key 的粒度,按账户、按股票、按股票加账户,状态设计和后续性能直接挂钩。
2.2 部署与资源规划:作业不是写完就能跑的
本地 IDE 调试跑起来的是简化版,生产环境通常是 YARN 或 Kubernetes 调度。不管哪种,作业提交到集群后都要跟调度器申请资源。并行度设计有个基本原则:Source 并行度最好等于 Kafka 分区数,别多也别少;窗口聚合算子并行度按数据量估算;Sink 并行度受下游存储连接能力限制,不是越大越好。
内存配置是新手的重灾区。TaskManager 的 total process memory 不只是堆内存,还包括 JVM overhead、网络缓冲、托管内存等。只看堆内存去调,作业跑起来没多久就可能 OOM。我一般建议用 RocksDB 做状态后端时,单 TaskManager 的托管内存至少预留 1~2GB,堆内存按业务复杂度另行分配。
Checkpoint 配置建议:证券场景 30~60 秒一次,最小间隔 10~20 秒。太频繁会导致大量快照和磁盘 IO,影响吞吐;太久则故障恢复时间过长,盘中挂掉再重启,几十秒的空窗期对实时指标影响很大。精确一次和至少一次的选择,取决于下游 Sink 是否幂等。ClickHouse 配合去重表引擎或幂等插入时相对容易,普通 MySQL 则需要数据库主键约束来保证幂等。
2.3 SpringBoot 整合 Flink:开发效率与作业提交的最佳平衡
很多入门同学搜“springboot整合flink”,第一反应是直接在一个 SpringBoot 应用里 new 一个 ExecutionEnvironment,从启动类开始跑。不是不能跑,而是只适合单机演示。生产环境中证券业务动辄几十个实时作业,不可能让每个作业都内嵌一个 JobManager。SpringBoot 在其中的定位更多是“控制面”:负责提交作业、管理参数、查看状态。
我采用过一种比较实用的方案:SpringBoot 提供 REST 接口,内部把作业 jar 和参数拼成 flink run 命令行提交:
public String submitJobToFlink(String jarPath, String mainClass, List<String> args) { List<String> command = new ArrayList<>(Arrays.asList( "flink", "run", "-d", "-m", "yarn-cluster", "-c", mainClass, jarPath )); command.addAll(args); ProcessBuilder pb = new ProcessBuilder(command); pb.redirectErrorStream(true); Process process = pb.start(); try (BufferedReader reader = new BufferedReader( new InputStreamReader(process.getInputStream()))) { String line; while ((line = reader.readLine()) != null) { log.info(line); } } return "job submitted"; }注意,生产环境里别直接依赖这个命令的输出文本判断成功失败,提交后要通过 Flink REST API 查询作业状态。另一种方案是直接调用 Flink REST API 上传 jar、触发作业,但要处理 multipart 上传、jobId 回执、状态轮询,代码量大一些。我推荐把提交控制模块独立成一个 SpringBoot 服务,实时作业代码打成独立 jar 做版本归档,这样提交、回滚、重启都清晰。
3. 核心链路实操:从 Kafka 到 ClickHouse 的实时聚合
3.1 一个可复现的行情实时聚合作业
先假设收到一条 Kafka 里的行情 tick 消息,JSON 格式长这样:
{"symbol":"600519","price":1500.5,"volume":20,"amount":30010.5,"ts":1700000000000}目标是每 10 秒输出每只股票的开盘、收盘、最高、最低、成交量。完整作业骨架如下:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(4); env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE); KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("kafka1:9092,kafka2:9092") .setTopics("stock-tick") .setGroupId("flink-market-analysis") .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStream<StockTick> tickStream = env.fromSource( source, WatermarkStrategy.noWatermarks(), "kafka-source" ).map(json -> JsonUtil.parse(json, StockTick.class)) .assignTimestampsAndWatermarks( WatermarkStrategy.<StockTick>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((tick, timestamp) -> tick.getTs()) ); SingleOutputStreamOperator<TickAggResult> aggStream = tickStream .filter(tick -> tick.getPrice() > 0) .keyBy(StockTick::getSymbol) .window(TumblingEventTimeWindows.of(Time.seconds(10))) .aggregate(new TickAggregateFunction(), new TickWindowProcessFunction());这里有个关键设计:aggregate 第一个参数是增量聚合函数,只维护少量状态;第二个参数是窗口处理函数,在窗口触发时仅调用一次,补上窗口起止时间。不要用 ProcessWindowFunction 去收集全窗口数据再遍历,行情数据量下必炸。
JdbcSink.sink( "INSERT INTO stock_agg(symbol, window_start, window_end, open, close, high, low, volume) " + "VALUES(?,?,?,?,?,?,?,?)", (ps, agg) -> { ps.setString(1, agg.getSymbol()); ps.setLong(2, agg.getWindowStart()); ps.setLong(3, agg.getWindowEnd()); ps.setBigDecimal(4, agg.getOpen()); ps.setBigDecimal(5, agg.getClose()); ps.setBigDecimal(6, agg.getHigh()); ps.setBigDecimal(7, agg.getLow()); ps.setLong(8, agg.getVolume()); }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchIntervalMs(3000) .withMaxRetries(3) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl("jdbc:clickhouse://clickhouse-host:8123/market") .withDriverName("com.clickhouse.jdbc.ClickHouseDriver") .withUsername("default") .withPassword("") .build() );写 Sink 之前,先用客户端工具验证一下 ClickHouse 驱动类名和连接串,不同驱动版本之间差异还挺大的。
3.2 用 Flink JDBC 连接器把 MySQL 同步到 ClickHouse
Flink 做 MySQL 到 ClickHouse 的实时同步,最常见的组合是 Flink CDC 监听 MySQL binlog 拿到变更事件,再经 JDBC Sink 写 ClickHouse。Flink CDC 连接器启动时会先做全量快照,再自动切到增量,不需要自己维护两条同步逻辑,非常省事。
Maven 依赖大致如下:
<dependency> <groupId>com.ververica</groupId> <artifactId>flink-connector-mysql-cdc</artifactId> <version>2.3.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-jdbc</artifactId> <version>1.16.0</version> </dependency>版本一定要跟 Flink 主版本匹配,官方文档里有兼容矩阵,不要直接写 latest。
核心代码骨架:
MySqlSource<String> mysqlSource = MySqlSource.<String>builder() .hostname("mysql-host") .port(3306) .databaseList("securities") .tableList("securities.order_record") .username("cdc_user") .password("****") .deserializer(new JsonDebeziumDeserializationSchema()) .build(); DataStreamSource<String> cdcStream = env.fromSource( mysqlSource, WatermarkStrategy.noWatermarks(), "mysql-cdc-source" );之后把 JSON 流解析成目标表结构,再通过 JdbcSink 写入 ClickHouse。需要强调一点:ClickHouse 不擅长高频点写,JDBC Sink 必须开批量,一次至少攒几百条再提交。高频少量写入会导致小文件碎片,后续查询性能和 merge 压力都会受影响。
还有一个容易踩的坑:如果业务源表有删除操作,CDC 流里会出现 DELETE 事件。ClickHouse 默认不支持单行删除,要么用 CollapsingMergeTree 或 ReplacingMergeTree 表引擎做逻辑删除,要么只同步允许追加的实时明细表。设计链路时就要把删除策略提前定好,否则上线后才发现数据对不上,返工成本很高。
3.3 JDBC 连接器参数调优与常见异常实录
JDBC 连接器是 Flink 生态里很常用的组件,但异常也不少。我整理了一份高频问题速查表:
| 异常现象 | 常见原因 | 处理办法 |
|---|---|---|
| ClassNotFoundException: com.mysql.cj.jdbc.Driver | 驱动 jar 没打进 fat jar,或依赖冲突被排除 | 用 maven-shade 打包,确认包含 mysql-connector-j;本地验证驱动类完整路径 |
| Connection is not available, request timed out | 连接池不够,高峰期并发请求超出池上限 | 增大 withBatchSize 减少交互次数,配置连接池参数,并行度不要盲目调大 |
| Batch flush failed, aborting | 目标表字段与结果类型不匹配,或超长文本、非法空值 | 检查出错行字段,统一用 setBigDecimal/setString,建表时做好默认值 |
| 数据重复、结果翻倍 | 重试导致重复写入,没有幂等机制 | 目标表加唯一键,使用 ReplacingMergeTree 或应用层去重 |
| 时间差了 8 小时 | ClickHouse 和 MySQL 会话时区不一致 | 连接串里显式指定 serverTimezone 和 use_time_zone,统一存 UTC |
JDBC 连接器的三个参数是关键:withBatchSize 不是越大越好,单次失败重试锁住的数据会更多;withBatchIntervalMs 控制缓存多久刷一次,实时要求高就设 1~3 秒;withMaxRetries 建议 3~5,超过后进入失败流程,让告警及时出现,而不是无限重试导致上游积压。
4. 生产环境问题排查与性能优化
4.1 背压:先看懂监控指标再动手优化
背压是流式系统最常见的问题。理解起来像河道:下游泄洪速度慢,上游水位就会上涨,整条链路的处理吞吐被拖垮。在 Flink Web UI 的 Back Pressure 标签页,能看到每个算子的背压比例。如果某个算子长时间处于高背压,就该定位了。
排查顺序很重要。先看 Sink 是否卡在下游存储,比如 ClickHouse 大批量插入时 merge 变慢,或者 MySQL 存在锁竞争;再看窗口算子 keyBy 后的数据分布;最后才怀疑 Source。别一上来就加并行度,加并行度往往会把压力放大到下游,反而更严重。
我在生产上遇到过一个很典型的背压场景:Kafka 分区 12 个,Flink Source 并行度调成 24,结果一半并行子任务在等数据,另一半空转,窗口算子状态局部倾斜更严重。后来把 Source 并行度改回 12,并给 keyBy 加随机前缀做两阶段聚合,整体吞吐直接翻倍。这个案例说明,背压优化不是无脑加资源,而是先找到真正的瓶颈。
4.2 Checkpoint 与状态恢复:实时任务的定海神针
证券场景强烈建议开启 Checkpoint。它存在的意义是作业挂掉后能从最近一次快照恢复,不会把累计指标全部丢掉。很多新手开了 Checkpoint 后,发现重启后数据重复或者恢复失败,问题多半出在几个地方。
第一,Sink 不是幂等的。Flink 的 EXACTLY_ONCE 是应用层语义,如果 Sink 本身重复写,最终结果表里还是会有重复数据。JDBC Sink 写入时,目标表要有合适的主键或者去重机制,ClickHouse 可以用 ReplacingMergeTree,MySQL 可以用唯一索引。
第二,Checkpoint 超时。默认超时时间可能偏短,Kafka Source 的 pending 数据很多或者 Sink 持有锁时,快照一直做不完。可以调大超时时间,同时设置最小间隔,让快照有喘息空间,但也不能太宽松,否则任务挂了半天才发现。
第三,RocksDB 状态大了以后,每次全量快照代价很高。建议开启增量快照:
RocksDBStateBackend rocksDB = new RocksDBStateBackend("hdfs://nameservice/flink/checkpoints", true); env.setStateBackend(rocksDB);Checkpoint 目录放在 HDFS,磁盘容量要提前规划。曾见过状态好几十 GB 的情况,如果快照目录爆了,整个作业会反复重启起不来。
4.3 数据延迟与乱序:如何保证结果不偏
实时计算的准确性很大程度上取决于如何处理乱序。forBoundedOutOfOrderness(Duration.ofSeconds(5))是最常用的 watermark 策略,意思是允许事件时间存在最多 5 秒的乱序,超过 5 秒的迟到数据不会进入当前窗口。证券行情经过网络和多级转发,延迟抖动一般在 1~2 秒,5 秒是个稳妥的初始值。
这里有个坑:如果某个 Kafka 分区长时间没有数据,watermark 不会推进,整个窗口永远不触发,结果就一直不出。建议读取 Kafka 时开启分区空闲检测:
WatermarkStrategy.<StockTick>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withIdleness(Duration.ofSeconds(30)) .withTimestampAssigner(...)再说迟到数据。allowedLateness 允许窗口触发后再等一段时间,晚到且还处于容忍范围内的数据,会再次触发计算并更新旧窗口结果。但注意,如果结果已经写进了 ClickHouse,后续更新就必须依赖同主键覆盖。不要盲目使用 allowedLateness,用了就要对上更新语义。
更推荐的做法是侧输出:把超出容忍范围的最后一批迟到数据单独收集起来,交给离线修正任务统一处理。很多团队习惯直接丢弃迟到数据,等发现告警漏掉一大片时,问题已经被掩盖很久了。
5. 证券实时分析场景的避坑清单
5.1 数据字段与精度:看似简单却最容易翻车
行情价格和成交金额,一开始很多人图省事用 double。但金融计算对精度极度敏感,用 double 累积 30010.5 和 30010.49 这种小数,最后聚合出来的成交额会跟业务系统对不上。结论很直接:金额、价格统一用 BigDecimal,能不用浮点就不用。Flink 默认序列化对 BigDecimal 不够高效,但正确性优先,状态量特别大时再想办法优化,比如用最小货币单位转成 long 存储,展示层再转换。
时间戳也经常出问题。某些上游系统给 13 位毫秒值,另一些给 10 位秒值,时区一混,窗口统计就偏 8 小时。我的处理方式是在接入层统一转成 UTC 毫秒 long,时间字段一律用 long 传递,不要在计算链路里反复用字符串格式化。
证券代码字段更乱。有的系统带交易所前缀,比如 SH600519,有的是纯六位代码 600519。做 keyBy 之前必须先把格式统一,否则同一只股票的行情流会被切到两个 key 上,聚合结果直接被拆成两半。要在入口处就把转换逻辑写死,不要散落在各条流里。
5.2 性能优化与数据倾斜
热门股和冷门股天然倾斜。某只热门股票的每分钟成交量可能是冷门股票的几百倍,按股票代码做 keyBy 后,一个算子实例会被热门股票压死,其他实例闲得没事干。解决办法是两阶段聚合:先给 key 加随机前缀做一轮局部聚合,再去掉前缀按真实 key 做最终聚合。这样热门 key 的数据先被分散了一次,压力均衡很多。前提是聚合函数满足交换律和结合律,像求和、最大最小、计数都可以;开盘价这种首条记录类的指标就不能这么做。
序列化性能也要注意。原始行情是 JSON,如果每次进入窗口都反复解析同一份数据,CPU 开销非常可观。尽量在 map 阶段一次性解析成 POJO,后面尽量直接操作对象,避免在 keyBy、窗口函数、Sink 里二次 parse。Flink 对 POJO 有内置序列化器,比反射驱动的 Kryo 高效不少。
JDBC Sink 的连接池也要针对数据量调整。默认连接池可能只有几个连接,数据量大时点写非常痛苦。一定要用批量提交,MySQL 连接串里打开 rewriteBatchedStatements=true,ClickHouse 的 HTTP 连接数按 Sink 并行度配置。这些细节往往就是背压的隐藏原因。
5.3 我总结的几条运维与开发经验
版本选择上,不要一上来就追最新。Flink 小版本之间 connector 兼容性、SQL 行为都可能变化,依赖版本跟 Flink 主版本不匹配时会直接启动失败。先在测试集群固定版本跑一段时间,所有作业验证通过再上生产。
依赖打包是个高频坑点。用统一的 maven-shade 插件配置,排除 flink-core、flink-runtime、log4j 这类会被运行时提供的包。提交作业时的 Main-Class 最好固定下来。我见过很多次提交时报 ProgramInvocationException,查了半天发现是类名大小写写错或者 jar 包漏了依赖,数据已经断了大半天。
监控告警要提前配好。作业里至少要输出关键指标,比如聚合记录数和当前 watermark,通过 Prometheus 上报到 Grafana。告警规则至少三条:作业 FAILED、checkpoint 连续失败三次、当前水位落后数据源超过三分钟。这三条规则能救绝大多数实时任务的命。
这个项目从单机 Demo 到生产集群,我最大的体会是:Flink 本身不难,难的是把证券业务的数据口径、模型和时序细节理清楚。不管是实时行情、分钟级聚合还是 MySQL 同步 ClickHouse,技术选型已经高度共识,真正拉开差距的是对数据质量的把控和对异常场景的应对。最近我又在尝试用 Flink SQL 替换一部分 DataStream 聚合逻辑,等跑顺了,再来写一篇链路对比。