ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink SQL + Kafka 实时统计实战:从建表、窗口聚合到踩坑排查

Flink SQL + Kafka 实时统计实战:从建表、窗口聚合到踩坑排查 做实时统计这件事我从 DataStream API 一路写到现在说实话早期被各种乱序、延迟、状态恢复的问题折磨得不轻。后来 Flink SQL 慢慢成熟尤其是和 Kafka 搭在一起用一条从消息队列到实时指标的链路变得非常干净Java 工程里起一个 Flink SQL 作业用 DDL 定义好 Kafka 来源表和结果去向表剩下的统计逻辑全交给 SQL。今天这篇文章就把我实际搭建这套程序的过程完整拆开从版本选型、建表 DDL、窗口聚合、结果 Sink到踩过的坑和排查思路全部按干货讲。适合正在学 Flink SQL 的 Java 后端也适合已经在用 Kafka 做数据管道、想快速产出实时统计结果的团队参考。为什么强调“Java 版”是因为很多同学搜到教程可能是 Python 脚本或者 Scala Demo真正落到 Java 工程里要配 connector、要处理依赖冲突、要设置状态后端细节全在工程代码里。我这里不会只贴一个 SQL 片段而是把整个作业从零到落地怎么组织讲清楚你跟着搭一遍基本就能迁移到自己业务上。1. 为什么我选择“Flink SQL Kafka”这套组合1.1 实时统计最朴素的链路先还原一个最常见的场景订单系统或者埋点日志系统每时每刻都在往 Kafka 里写消息。业务方想实时知道的往往就是这几类问题现在每分钟有多少笔订单、累计成交额是多少、哪个品类卖得最好、今天有多少活跃用户。这些指标如果靠传统 Java 后端做定时任务去数据库里扫数据量一大根本扛不住而且统计口径分散在业务代码里改一个指标要重新发版。所以需要一个独立于业务系统之外的实时计算层Kafka 做数据通道Flink SQL 做计算引擎。这条链路有固定套路Kafka topic 进Flink 作业算结果写回 Kafka 或者 MySQL/ES。我长期用这套方案的原因很朴素。第一Kafka 作为消息缓冲层可以削峰业务端只管发送不关心下游怎么算第二Flink SQL 的语义足够直观一个窗口聚合就是一条 CREATE TABLE 加一条 INSERT INTO SELECT不需要像手写 DataStream 那样处理底层各类细节第三状态和容错是引擎自带的作业重启后能从上一次 checkpoint 恢复不用自己维护中间结果。1.2 Flink SQL 相比 DataStream API 的核心优势早期很多团队排斥 Flink SQL觉得灵活性不够只能写 DataStream API 才叫“真实时计算”。实际上实时统计类的需求有非常明显的模式化特征窗口聚合、多维计数、去重、TopN、累计值这些用 SQL 表达几乎是标准答案。对比维度Flink SQLDataStream API开发速度声明式 SQL指标改动快需要写 Java 逻辑迭代慢窗口/水印DDL 里声明即可手动 assignTimestampsAndWatermarks状态管理引擎托管可配置 TTL自己控制 ValueState/MapState维护成本一行 SQL 能看懂逻辑藏在代码里接手成本高适合场景多维度统计、窗口聚合、TopN复杂事件处理、自定义 UDF、多流 join但也不是完全替代。当业务需要多流按业务 ID 匹配、复杂状态机、或者要实时调用外部服务做特征加工时还是要回到 DataStream API。以我自己的经验来看SQL 覆盖了简单统计和报表类 80% 的场景剩下 20% 需求写成 UDF 嵌入 SQL 也不难。我的建议是先 SQLSQL 表达不了再加 UDF最后才考虑 DataStream不要一上来就用低级 API 堆逻辑。2. 环境准备与工程项目搭建2.1 版本选型Flink 生态的版本兼容是一件很头疼的事版本选不对后面连 ClassNotFound 都算轻的还可能遇到 connector 和引擎之间方法签名对不上的情况。我这里给出一套我实际跑通的组合JDK8 或 11 都行我建议 JDK 11长期维护GC 表现也好一些Flink1.17.2这是目前线上用得比较稳的版本新特性不少社区踩坑资料也全Kafka 3.4.02.8 以上的版本基本都兼容Flink SQL Kafka Connectorflink-sql-connector-kafka 1.17.2注意是带 sql 前缀的 fat jarJDBC Connectorflink-connector-jdbc 3.1.2-1.17用于写 MySQL。Java 环境变量配置属于最基础的前置工作这里默认你已经做好了。如果你还没配好 JAVA_HOME先去把 JDK 装好否则起 Flink 作业第一步就会卡在环境上。2.2 Maven 依赖怎么写工程用 Maven 管理依赖pom.xml 里核心依赖如下properties flink.version1.17.2/flink.version /properties dependencies !-- Flink Table API -- dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-planner-loader/artifactId version${flink.version}/version scopeprovided/scope /dependency !-- Flink SQL Kafka Connector需要打包进作业 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-sql-connector-kafka/artifactId version1.17.2/version /dependency !-- JDBC Sink -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version3.1.2-1.17/version /dependency !-- MySQL 驱动 -- dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency !-- RocksDB 状态后端建议加上 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-statebackend-rocksdb/artifactId version${flink.version}/version /dependency /dependencies注意几个细节。flink-table-api-java-bridge 和 flink-table-planner-loader 的 scope 是 provided因为 Flink 集群本身已经有这些类打进去反而容易和集群冲突。flink-sql-connector-kafka 必须打进作业 jar因为它不是 Flink 发行版自带的需要单独准备。flink-connector-jdbc 同理也需要打进去。打包用 maven-shade-plugin别用 assemblyshade 会正确重定位部分依赖类减少冲突。如果你的作业是在 Flink SQL Client 里跑也可以不用 Java 工程直接把 connector jar 放到 Flink 的 lib 目录然后写纯 SQL 脚本。但 Java 工程的好处是便于接入公司内部的配置中心、统一日志、监控上报以及写自定义 UDF。2.3 连接 Kafka 集群的前置检查很多同学一上来直接跑 Flink 作业结果报错然后对着日志一头雾水。我的习惯是先用命令行确认 Kafka 本身是通的再启动 Flink 作业这样可以快速隔离问题。先建 topic 并手动发一条消息验证kafka-topics.sh --bootstrap-server kafka-1:9092 --create --topic order_topic --partitions 6 --replication-factor 1 kafka-console-producer.sh --bootstrap-server kafka-1:9092 --topic order_topic {order_id:o_10001,user_id:u_8888,product_id:p_233,category:数码,amount:1299.00,ts:2024-06-01 10:15:30}然后另开一个终端消费确认消息能读到kafka-console-consumer.sh --bootstrap-server kafka-1:9092 --topic order_topic --from-beginning如果生产和消费都正常说明 Kafka 集群对外服务没问题问题只可能出在 Flink 作业的配置上。如果这里就报错先解决 Kafka 侧的网络、权限、topic 分区等问题再往下走。“error while fetching metadata”这类报错最常见原因有三个bootstrap.servers 配置不对、advertised.listeners 和客户端不在同一个网段、安全组拦截了 9092 端口。排查顺序是先 telnet 看端口通不通再看 advertised.listeners 配置的地址从客户端能否访问最后确认 topic 是否存在。尤其在 Docker 部署 Kafka 的场景容器内部监听地址和宿主机外部访问地址经常不一致这是超高频踩坑点。3. 用 Flink SQL 把 Kafka 数据读进来3.1 先定数据模型与消息格式Flink SQL 建表前一定要先想清楚 Kafka 里的消息长什么样。以订单为例我常用的 JSON 消息格式是这样的{ order_id: o_10001, user_id: u_8888, product_id: p_233, category: 数码, amount: 1299.00, ts: 2024-06-01 10:15:30 }字段类型上注意几点。amount 建议用 DECIMAL(10, 2)不要用 DOUBLE避免金额精度问题。ts 时间字段映射成 TIMESTAMP(3)毫秒精度。user_id、order_id 这类 ID 字段用 STRING别用 BIGINT因为很多 ID 在前端会被处理成字符串强行转 BIGINT 遇到脏数据会直接导致作业失败。Kafka 消息格式用 JSON 还是 CSV取决于上游系统。JSON 可读性好、字段扩展方便缺点是消息体积比 CSV 大吞吐会受一些影响。如果对性能要求极高可以考虑 Avro 加 Schema Registry但维护成本会上升小团队没必要一上来就上 Avro。3.2 创建 Kafka Source 表在 Flink SQL 里读 Kafka本质就是定义一张带 connector 属性的表。DDL 如下CREATE TABLE order_source ( order_id STRING, user_id STRING, product_id STRING, category STRING, amount DECIMAL(10, 2), ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic order_topic, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id flink-order-stat-group, scan.startup.mode earliest-offset, format json );这个 DDL 里最值得展开的是 WATERMARK 那一行。实时数据从业务系统发到 Kafka再被 Flink 消费过程中不可避免会有网络抖动、业务系统处理延迟导致消息到达 Flink 的顺序和实际发生时间不一致。如果不处理乱序直接按事件的 ts 去开窗口做聚合结果会偏差很大。水印就是用来告诉 Flink“我最多容忍 5 秒的乱序超过这个迟到范围的数据窗口就不等了。”WATERMARK FOR ts AS ts - INTERVAL 5 SECOND 的具体含义是当前水印时间 已收到消息中最大的 ts 减去 5 秒。窗口触发条件是水印时间 窗口结束时间。所以一个 10:15:00 到 10:16:00 的窗口要等到水印推进到 10:16:00 才会触发计算也就是要收到一个 ts 至少为 10:16:05 的消息。如果你不需要精确的事件时间统计只想看个大概可以完全不用事件时间建表时不声明 ts 和 WATERMARK窗口直接用处理时间 PROCETIME()。但做实时统计尤其是涉及金额、订单量这类要和业务对账的指标我强烈建议用事件时间加水印结果才可信。3.3 消费位点与并行度注意点scan.startup.mode 决定 Flink 作业启动时从 Kafka 的哪个位置开始消费常用取值有这几个earliest-offset从头开始消费适合首次上线想回补历史数据latest-offset只消费作业启动后新到的数据适合不关心历史、只要增量timestamp从指定时间戳之后开始消费适合想从某个业务时间点补数据specific-offsets从指定分区的指定 offset 开始消费一般用得少。线上首次上线一个统计作业如果 topic 里已经有堆积数据我又想快速看到结果一般用 timestamp指定到当前时间往前推 5 分钟既不会重复消费太多历史也不会漏掉启动瞬间产生的消息。如果数据量小、可以接受重新计算历史用 earliest-offset 最省事。并行度方面Flink SQL 作业的 source 并行度理论上可以大于 Kafka 分区数但没意义多的并行度只会空闲。合理的做法是并行度取分区数的整数倍比如 topic 6 个分区作业并行度设 6 或者 12调度最均衡。如果你用 Flink SQL Client 直接跑可以用 SET parallelism.default 6 来控制。4. 实时统计逻辑怎么用 SQL 表达4.1 最常用的几个实时指标模板实时统计的需求翻来覆去就那么几类我把常用的指标模板整理出来对应到 SQL 上就是固定套路每分钟订单数和 GMV用 TUMBLE 滚动窗口按分钟切分窗口结束触发输出近 1 小时每 10 分钟更新的滚动统计用 HOP 滑动窗口窗口长度 1 小时滑动步长 10 分钟全量累计 PV/UV用 COUNT 和 COUNT(DISTINCT)但要小心状态膨胀品类销量 TopN窗口内先聚合再用 ROW_NUMBER 分组排序会话级统计用 SESSION 会话窗口按用户空闲时间切分会话。刚开始做实时统计建议先从 TUMBLE 滚动窗口入手因为它最简单直观窗口之间互不重叠理解成本最低。跑通了再加 HOP 和 SESSION理解窗口的触发机制后其他窗口只是切分规则不同。4.2 事件时间窗口聚合完整示例下面这个例子是典型的“每分钟统计每个品类的订单数和 GMV”并且直接把结果写入 MySQL。先建结果表CREATE TABLE order_stats_sink ( window_start TIMESTAMP(3), window_end TIMESTAMP(3), category STRING, order_cnt BIGINT, gmv DECIMAL(14, 2), PRIMARY KEY (window_start, window_end, category) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/stats, table-name order_stats, username root, password 123456 );然后提交计算逻辑INSERT INTO order_stats_sink SELECT TUMBLE_START(ts, INTERVAL 1 MINUTE) AS window_start, TUMBLE_END(ts, INTERVAL 1 MINUTE) AS window_end, category, COUNT(*) AS order_cnt, SUM(amount) AS gmv FROM order_source GROUP BY TUMBLE(ts, INTERVAL 1 MINUTE), category;这里有一个非常关键的细节GROUP BY 后面跟着窗口函数和 category说明每个品类在每个分钟窗口内都会产出一条结果。如果 category 的取值很多比如几百个品类一分钟就会产生几百条结果下游存储压力不小设计结果表时要有这个预期。TUMBLE_START 和 TUMBLE_END 是窗口开始和结束时间必须显式查询出来否则下游看到的只有聚合值不知道对应哪个时间窗口。用 PRIMARY KEY NOT ENFORCED 声明主键是为了让 JDBC sink 在写入时能按主键做更新而不是每条都追加。这个声明不会在 Flink 侧强制校验只影响 sink 的写入语义。4.3 窗口 TopN 和累计去重订单量和 GMV 只是基础指标业务方很快会追问“哪个品类卖得最好”。这就是 TopN 需求。Flink SQL 的窗口 TopN 写法如下SELECT window_start, window_end, category, gmv FROM ( SELECT window_start, window_end, category, gmv, ROW_NUMBER() OVER ( PARTITION BY window_start, window_end ORDER BY gmv DESC ) AS rn FROM ( SELECT TUMBLE_START(ts, INTERVAL 5 MINUTE) AS window_start, TUMBLE_END(ts, INTERVAL 5 MINUTE) AS window_end, category, SUM(amount) AS gmv FROM order_source GROUP BY TUMBLE(ts, INTERVAL 5 MINUTE), category ) ) WHERE rn 10;注意这里用的是窗口 TopN不是普通 TopN。普通 TopN 在无界流上会一直维护全局排名数据量一大状态就爆炸。窗口 TopN 按 window_start 和 window_end 分区每个窗口只保留前 10 名窗口结束状态就可以清理逻辑和性能都可控。再讲 UV 类指标。最直接的写法是 COUNT(DISTINCT user_id)但 Flink SQL 默认的 COUNT(DISTINCT) 是精确去重状态里会保存每一个出现过的 user_id用户量一大状态后端会撑不住。如果 UV 量级在百万以下精确去重可以接受如果到千万级以上建议用近似去重方案比如自定义 HyperLogLog 的 UDAF或者把去重窗口缩短比如按小时而不是按天去重。实时报表的 UV 本来就允许一定误差1% 以内的偏差业务方几乎感知不到。5. 结果写到哪里Sink 设计与落地5.1 写回 Kafka 做下游复用实时统计结果最常见的去向之一是再写回 Kafka。这样下游的 SpringBoot 服务、ES、数据仓库都可以各自消费结果 topic互不影响。这个思路和“实时数仓分层”很像中间结果层先落到 Kafka下游再按需取用。写回 Kafka 的 DDL 非常简单CREATE TABLE result_kafka ( window_start TIMESTAMP(3), category STRING, order_cnt BIGINT, gmv DECIMAL(14, 2) ) WITH ( connector kafka, topic order_stats_topic, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, format json ); INSERT INTO result_kafka SELECT TUMBLE_START(ts, INTERVAL 1 MINUTE), category, COUNT(*), SUM(amount) FROM order_source GROUP BY TUMBLE(ts, INTERVAL 1 MINUTE), category;消费端如果是 SpringBoot配置好 Kafka 的 consumer 就能拿到 JSON 格式的统计结果直接推页面或者做阈值告警都行。这种链路我实际用下来最大的好处是解耦统计结果和业务系统之间没有强依赖下游哪怕挂了结果还在 Kafka 里躺着恢复后重新消费就行。5.2 JDBC 写入 MySQL 的参数推荐如果统计结果要直接对接到报表系统或者管理后台写到 MySQL 是最省事的方式。除了 4.2 里的 JDBC 建表 DDL有几个参数我建议你重点关注sink.buffer-flush.max-rows默认 100建议调大到 1000减少 flush 次数sink.buffer-flush.interval默认 1s建议设成 2s让数据攒一攒再批量写sink.parallelism不要设置太高2 到 4 就够了太高会给 MySQL 造成不必要的连接和写入压力。JDBC sink 默认的写入语义是 At Least Once作业重启时可能出现重复写入。要保证下游数据最终一致一个简单可靠的做法是 MySQL 结果表设计联合唯一键比如把 window_start、window_end、category 三个字段建唯一索引然后底层用 INSERT ON DUPLICATE KEY UPDATE 的语义去更新。Flink JDBC connector 本身不直接支持这种 upsert 语法所以很多团队会把结果先写到 Kafka再用一个轻量消费者异步写入 MySQL这是另一个话题了但思路可以记下来。5.3 Checkpoint 和状态后端不配置等于白做Flink 的容错全部依赖 checkpoint这句话我在多个场合强调过。如果 checkpoint 没开作业一重启状态全丢Kafka offset 也要从配置的起始位置重新消费结果就是重复计算、数据错乱。实时统计作业checkpoint 是必须配置的。Java 代码里配置如下StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 每 60 秒触发一次 checkpoint env.enableCheckpointing(60_000); // 精确一次语义 env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 两次 checkpoint 之间最少间隔 30 秒避免频繁快照 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000); // checkpoint 存储位置生产环境建议 HDFS env.getCheckpointConfig().setCheckpointStorage(file:///opt/flink/checkpoints); env.setStateBackend(new EmbeddedRocksDBStateBackend()); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env);RocksDB 状态后端是我上线 Flink SQL 作业的默认选择。相比 HashMapStateBackendRocksDB 把状态存储从 JVM 堆内挪到了堆外磁盘状态量再大也不会 OOM只是读写性能比纯内存差一些。实时统计作业的状态通常不小用 RocksDB 更稳。SQL 侧还可以设置状态 TTL把过期状态清理掉SET table.exec.state.ttl 2h;这个配置对 GROUP BY 和 JOIN 都生效。尤其对大窗口、大 key 数量的统计TTL 能明显控制状态大小。但注意别设置太短TTL 太短会导致窗口还没结束状态已经被清理统计结果瞬间变错。我一般按窗口长度的 3 到 5 倍来设置。6. 我踩过的坑和排查实录6.1 时间字段与时区差 8 小时这个坑几乎所有用 Flink SQL 的人都躲不过。业务库里的时间字段通常是北京时间但 Flink 的 TIMESTAMP 类型默认按 UTC 处理窗口边界如果按 UTC 计算你看到的统计分钟窗口会比真实时间慢 8 小时。解决办法是把时间字段定义成 TIMESTAMP_LTZ(3)并设置 Flink 的本地时区SET table.local-time-zone Asia/Shanghai;建表时这样声明CREATE TABLE order_source ( ... ts TIMESTAMP_LTZ(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH (...);TIMESTAMP 和 TIMESTAMP_LTZ 的区别简单记就是TIMESTAMP 不带时区看到什么就是什么TIMESTAMP_LTZ 带时区语义Flink 会根据本地时区转换后参与窗口计算。如果 Kafka 里存的是 epoch 毫秒数用 TIMESTAMP_LTZ 会非常合适。存的是格式化字符串时间建议先用 TO_TIMESTAMP_LTZ 转成 TIMESTAMP_LTZ 再参与计算。6.2 Kafka 消费延迟高怎么定位消费延迟高是实时统计最常见的运维问题。表现就是作业还在跑但结果指标明显落后于当前时间。定位延迟第一件事不是看 Flink 日志而是看 Kafka 消费组的 LAGkafka-consumer-groups.sh --bootstrap-server kafka-1:9092 --group flink-order-stat-group --describe输出里能看到每个分区的 LOG-END-OFFSET 和 CURRENT-OFFSET两者差值就是积压量。如果 LAG 持续增长说明消费速度跟不上生产速度。常见原因按概率排序sink 写入 MySQL 太慢、并行度不足、状态太大导致 checkpoint 耗时变长、GC 频繁。排查顺序先看下游再看并行度最后才看状态。曾经遇到过一次统计作业 LAG 飙升查到最后是 MySQL 结果表的一个索引失效导致 sink 单条写入越来越慢。去掉索引并调大 buffer-flush 参数后LAG 立刻降下来。实时作业的下游存储性能往往是最先出问题的环节。6.3 状态无限膨胀导致作业 OOMCOUNT(DISTINCT)、大窗口聚合、普通 TopN 都会让状态越来越大。我踩过一次 UV 状态把 RocksDB 撑到几十 GB 的坑起因是对用户 ID 做了精确去重且没有设置状态 TTL状态只增不减。应对策略有三层。第一层合理使用窗口化统计让状态随着窗口结束自动清理第二层设置 table.exec.state.ttl把超过窗口周期的旧状态清掉第三层对去重类指标改用近似算法比如 HyperLogLog牺牲少量精度换状态量级的下降。这里特别提醒一点不要盲目调大 JVM 堆内存来扛状态。状态用 RocksDB 落到磁盘后内存不够时可以用磁盘存状态但磁盘也不是无限的最终还是要从业务逻辑上减少状态量。实时统计的常态是“能用窗口表达就别用全局累加”。6.4 Kafka 连接报错速查表整理了我在实操中遇到过的几个高频 Kafka 连接问题放在一起方便排查报错现象可能原因排查方向error while fetching metadata with correlation idbootstrap.servers 不可达 / advertised.listeners 地址错误telnet 端口、检查 listeners 配置LEADER_NOT_AVAILABLEtopic 刚创建分区 leader 选举未完成稍等重试确认 broker 状态正常Connection timed out网络不通、安全组拦截、跨网段访问telnet 端口、检查防火墙策略Disconnected while connectingKafka 端启用了 SSL/SASL客户端未配置检查是否启用安全认证补认证参数SerializationExceptionJSON 格式解析失败字段类型不匹配用 console-consumer 看原始消息比对字段类型遇到 Kafka 连接问题先确认基础网络再往上查配置。很多人一上来就翻 Flink 日志其实 Flink 的报错信息已经明确写了“Disconnected”或者“Timeout”根因往往就是网络或者 listeners 配置和 Flink 本身没关系。做实时统计这么久我个人的体会是Flink SQL 把计算逻辑的门槛降得很低难点已经从“怎么写聚合”转移到了“怎么把链路的每个环节配置正确”。从 Kafka 消费的位点、水印的设定、窗口的选择到 sink 的幂等、checkpoint 的恢复每一环都要对背后的原理有清晰认知不然作业跑起来是一回事跑得稳不稳定、数据准不准确又是另一回事。最后再分享一个落地建议。第一次上手不要追求指标多先拿一个订单 topic 做分钟级 GMV 统计把这个链路完整跑通再逐步加窗口 TopN、UV、多维对比。同时把 checkppoint、状态后端、日志监控这些基础设施在第一个作业就配好后面加作业就只是加 DDL 和 SQL 模板的事。这套东西一旦跑顺后面新需求基本就是流水线作业了。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进