ARTICLE DETAIL

资讯详情

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

Flink在证券实时数据分析中的架构设计与生产实践

Flink在证券实时数据分析中的架构设计与生产实践 做证券行业的实时数据分析最折磨人的不是写不出 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, ListString args) { ListString 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); KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka1:9092,kafka2:9092) .setTopics(stock-tick) .setGroupId(flink-market-analysis) .setStartingOffsets(OffsetsInitializer.latest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamStockTick tickStream env.fromSource( source, WatermarkStrategy.noWatermarks(), kafka-source ).map(json - JsonUtil.parse(json, StockTick.class)) .assignTimestampsAndWatermarks( WatermarkStrategy.StockTickforBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((tick, timestamp) - tick.getTs()) ); SingleOutputStreamOperatorTickAggResult 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 同步到 ClickHouseFlink 做 MySQL 到 ClickHouse 的实时同步最常见的组合是 Flink CDC 监听 MySQL binlog 拿到变更事件再经 JDBC Sink 写 ClickHouse。Flink CDC 连接器启动时会先做全量快照再自动切到增量不需要自己维护两条同步逻辑非常省事。Maven 依赖大致如下dependency groupIdcom.ververica/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.3.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version1.16.0/version /dependency版本一定要跟 Flink 主版本匹配官方文档里有兼容矩阵不要直接写 latest。核心代码骨架MySqlSourceString mysqlSource MySqlSource.Stringbuilder() .hostname(mysql-host) .port(3306) .databaseList(securities) .tableList(securities.order_record) .username(cdc_user) .password(****) .deserializer(new JsonDebeziumDeserializationSchema()) .build(); DataStreamSourceString 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统一存 UTCJDBC 连接器的三个参数是关键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 可以用 ReplacingMergeTreeMySQL 可以用唯一索引。第二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.StockTickforBoundedOutOfOrderness(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 连接串里打开 rewriteBatchedStatementstrueClickHouse 的 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 聚合逻辑等跑顺了再来写一篇链路对比。
返回列表