简介这份资源面向大数据开发初学者与需要搭建实时数仓的工程师聚焦Flink实时消费Kafka数据、按定时或数量条件批量聚合后写入MySQL的完整实现。压缩包共9个文件以4个Java源码为核心配合2个SQL建表脚本、1个pom.xml依赖配置以及Kafka与Zookeeper的tgz、gz安装包整体约67.84MB便于直接搭建可运行的实时处理环境。源码演示了FlinkKafkaConsumer实时摄入、时间窗口与计数窗口两种触发策略以及通过JDBC或Table API将聚合结果持久化到关系库的写法同时附带Kafka集群所需的Zookeeper组件省去单独寻找版本匹配安装包的麻烦。已有3418人学习下载适合对照代码理解流处理链路、快速复现实验并迁移到自身业务场景。1. Flink 读 Kafka 攒批写 MySQL为什么“定时按量”双触发的 Sink 更靠谱Kafka 里的数据一条一条往 MySQL 写是很多人上手 Flink 的第一个版本也是最快翻车的一版。每条消息触发一次 JDBC 连接和一次INSERTQPS 一上来连接池先崩MySQL 的 binlog 和磁盘 IO 跟着告警端到端延迟反而比攒批更高。这个标题要解决的就是这件事用 Flink 实时消费 Kafka在算子内部做批量聚合按“攒够 N 条”或“等够 T 秒”两个条件谁先满足谁触发再一次性写进 MySQL。它适合正在做实时数仓 ODS/DWD 落库、埋点明细入库、订单流水同步的工程师也适合被“一条一写”坑过、想找一个能直接复现的攒批方案的人。核心不是 Flink 会不会连 Kafka而是批量聚合的触发时机怎么定、状态怎么管、失败怎么不丢不重这三件事决定了这套链路能不能上生产。2. 攒批 Sink 的三种实现路线先选对再动手2.1 为什么不用 Flink 自带的 JDBC Connector 直接写Flink 官方生态里有 JDBC Connector配置里能设sink.buffer-flush.max-rows和sink.buffer-flush.interval看起来正好是“按数量定时”。但实际用下来有几个绕不开的点它的攒批是在 Sink 内部做的批与批之间没有业务语义你没法在写库前对这批数据做二次聚合、去重或字段加工一旦 MySQL 报错整批重试的粒度是 Connector 内部决定的你想按业务主键做幂等 upsert 很别扭再加上不同 Flink 版本里这个 Connector 的包名和参数名改过几轮升级时容易踩到flink的jdbc连接器异常这类问题。所以只要你的场景里带一点业务加工或者对“这批数据到底写了什么”有可观测性要求自己写一个RichSinkFunction或ProcessFunction更可控。2.2 三条路线对比RichSinkFunction、ProcessFunction、自定义 Flink Sink路线攒批位置定时能力状态与容错适用场景RichSinkFunctioninvoke 内维护 List靠独立线程或检查时间戳需自己接 Checkpoint逻辑简单、批大小固定KeyedProcessFunction 定时器算子状态原生 event timer / processing timer状态后端托管天然对齐 Checkpoint需要按 key 聚合、按事件时间触发自定义 Sink实现 Sink 接口Sink 内部依赖 Sink 生命周期与 Flink 事务机制耦合深需要两阶段提交、精确一次我一般会选第二条用KeyedProcessFunction或ProcessFunction维护一个ListState注册一个 processing time 定时器做兜底每来一条数据判断数量是否到阈值到了就立刻 flush 并清掉定时器。这样“按数量”是主触发“定时”是保底触发两个条件互斥且都能触发写库逻辑清晰状态也能被 Checkpoint 托管。2.3 最小可跑骨架Kafka Source 到攒批算子先搭一个能跑通的最小骨架把 Kafka 消费和攒批逻辑串起来。下面这段是核心结构省略了具体的 MySQL 写入细节下一章展开。// 攒批算子的核心骨架数量触发 定时兜底 public class BatchAggregateFunction extends KeyedProcessFunctionString, String, ListString { // 批大小阈值攒够这么多条就立刻 flush private final int batchSize; // 定时兜底间隔单位毫秒 private final long intervalMs; // 用 ListState 存当前批次交给状态后端托管 private ListStateString bufferState; // 记录已注册的定时器时间戳避免重复注册 private ValueStateLong timerState; public BatchAggregateFunction(int batchSize, long intervalMs) { this.batchSize batchSize; this.intervalMs intervalMs; } Override public void open(Configuration parameters) { ListStateDescriptorString bufferDesc new ListStateDescriptor(buffer, String.class); bufferState getRuntimeContext().getListState(bufferDesc); timerState getRuntimeContext().getState( new ValueStateDescriptor(timer, Long.class)); } Override public void processElement(String value, Context ctx, CollectorListString out) throws Exception { bufferState.add(value); // 第一次进来时注册一个兜底定时器 if (timerState.value() null) { long triggerAt ctx.timerService().currentProcessingTime() intervalMs; ctx.timerService().registerProcessingTimeTimer(triggerAt); timerState.update(triggerAt); } // 数量触发攒够 batchSize 立即输出并清理 int count 0; for (String ignored : bufferState.get()) { count; } if (count batchSize) { flush(out, ctx); } } Override public void onTimer(long timestamp, OnTimerContext ctx, CollectorListString out) throws Exception { // 定时触发不管攒了多少都刷出去 flush(out, ctx); } private void flush(CollectorListString out, Context ctx) throws Exception { ListString batch new ArrayList(); for (String s : bufferState.get()) { batch.add(s); } if (!batch.isEmpty()) { out.collect(batch); } bufferState.clear(); // 清掉定时器状态下次数据进来重新注册 Long t timerState.value(); if (t ! null) { ctx.timerService().deleteProcessingTimeTimer(t); timerState.clear(); } } }逻辑说明processElement每来一条数据先入ListState然后判断当前批是否达到batchSize达到就调用flush输出整批并清空状态同时用timerState保证每个批次只注册一个定时器onTimer到点后无论攒了多少都强制 flush。参数说明batchSize决定吞吐和延迟的平衡点intervalMs决定最坏情况下一条数据要等多久才落库。这两个值不是拍脑袋定的下一章讲怎么算。3. 把参数调对批大小、定时器与 MySQL 写入的配合3.1 batchSize 和 intervalMs 怎么算不是越大越好批大小和定时间隔是一对互相牵制的参数。批越大MySQL 的INSERT次数越少吞吐越高但单批占用的堆内存越大且一旦这批写失败重试的数据量也越大间隔越长延迟越高但能保证低峰期数据不会一直卡在内存里。我的经验算法是先看 MySQL 单次批量写入的舒适区一般单条INSERT ... VALUES (...),(...),...拼到 500 到 2000 行比较稳再大容易撞max_allowed_packet然后看业务能接受的最大延迟埋点类通常 5 到 10 秒订单类 1 到 3 秒。于是batchSize取 500 到 1000intervalMs取 2000 到 5000先跑起来看 Kafka lag 和 MySQL 写入耗时再微调。// 参数配置建议从保守值起步压测后再调 int batchSize 500; // 单批最多 500 条 long intervalMs 3000L; // 最多等 3 秒必刷 // 如果 Kafka 峰值 QPS 是 20003 秒最多攒 6000 条 // 但 batchSize 是 500所以实际是数量触发为主定时器兜底低峰这里有个容易忽略的点如果峰值 QPS 很高数量触发会一直生效定时器几乎不会触发如果低峰 QPS 很低定时器才是主力。所以两个参数要按“峰值靠数量、低峰靠定时”的思路配而不是只调一个。3.2 用 JDBC 批量写入 MySQLrewriteBatchedStatements 必须开攒好的一批数据最终要写进 MySQL。用 JDBC 的addBatchexecuteBatch是标准做法但很多人不知道 MySQL 驱动默认不会把多条INSERT合并成一条导致“批量”其实是假批量网络往返次数没减少。必须在连接串上开rewriteBatchedStatementstrue驱动才会把批内语句重写成一条多值INSERT。// MySQL 批量写入连接串和 PreparedStatement 都要配对 String url jdbc:mysql://127.0.0.1:3306/demo ?rewriteBatchedStatementstrue // 关键让驱动合并多值 INSERT useServerPrepStmtstrue cachePrepStmtstrue connectTimeout3000 socketTimeout10000; String sql INSERT INTO order_flow (order_id, user_id, amount, ts) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE amount VALUES(amount), ts VALUES(ts); try (Connection conn DriverManager.getConnection(url, user, pwd); PreparedStatement ps conn.prepareStatement(sql)) { conn.setAutoCommit(false); for (String row : batch) { // 假设 row 已解析成字段这里做占位 ps.setLong(1, orderId); ps.setLong(2, userId); ps.setBigDecimal(3, amount); ps.setLong(4, ts); ps.addBatch(); } ps.executeBatch(); // 一次网络往返写入整批 conn.commit(); }逻辑说明rewriteBatchedStatementstrue是 MySQL 批量写入的性能开关不开的话executeBatch会退化成逐条发送ON DUPLICATE KEY UPDATE让重试时按主键幂等避免重复写。参数说明useServerPrepStmts和cachePrepStmts配合使用能减少预编译开销socketTimeout要设得比单批写入耗时长否则大批量写入时容易超时断开。3.3 连接池与重试别让一次 MySQL 抖动拖垮整条链路攒批写入最怕的是 MySQL 短暂不可用。如果 flush 时直接抛异常Flink 任务会失败重启从 Checkpoint 恢复后重新消费虽然不丢但会重复。更稳的做法是在 Sink 内部做有限次重试配合连接池管理连接。// 带重试的批量写入最多重试 3 次指数退避 int maxRetry 3; for (int i 0; i maxRetry; i) { try { writeBatch(batch); // 内部用 HikariCP 拿连接 break; } catch (SQLException e) { if (i maxRetry - 1) { throw e; // 重试耗尽交给 Flink 容错 } Thread.sleep((long) Math.pow(2, i) * 200); // 200ms, 400ms, 800ms } }逻辑说明重试只针对可恢复的SQLException重试间隔指数退避避免 MySQL 刚恢复就被打满。参数说明maxRetry不宜过大3 次足够覆盖瞬时抖动再长会拖高延迟连接池用 HikariCPmaximumPoolSize设成并行度乘以 2 左右别设太大把 MySQL 连接数占满。4. 避坑与排查攒批写 MySQL 最容易翻车的 5 个点4.1 现象任务跑一段时间后 Kafka lag 持续上涨但 MySQL 写入没报错原因通常是攒批算子的并行度和 Kafka 分区数不匹配或者batchSize设得太大导致单批处理时间过长算子成为瓶颈。解决先看 Flink Web UI 里算子的 busy 和 backpressure 指标如果 backpressure 高把batchSize调小到 200 到 500或者提高算子并行度对齐 Kafka 分区数。另外检查 MySQL 侧是不是有慢查询executeBatch单批耗时超过intervalMs时定时器会不断堆积。4.2 现象MySQL 里出现重复数据主键冲突或行数比 Kafka 多原因是任务失败重启后从 Checkpoint 恢复上一批已经写入 MySQL 但 Checkpoint 还没完成恢复后重新消费这段数据。解决写入语句必须带幂等语义用ON DUPLICATE KEY UPDATE或INSERT IGNORE主键选业务唯一键而不是自增 ID。如果业务允许也可以开启 Flink 的 Checkpoint 对齐和两阶段提交但实现复杂度高多数场景用幂等写入更划算。4.3 现象定时器不触发低峰期数据一直不落库原因是timerState在 flush 后没有正确清理或者注册定时器时用了 event time 但数据里没有水位线推进。解决确认用的是currentProcessingTime注册 processing time 定时器它不依赖数据时间检查flush里是否调用了deleteProcessingTimeTimer并timerState.clear()否则下一批数据进来时timerState.value()不为 null不会再注册新定时器。4.4 现象MySQL 报Packet for query is too large原因是单批拼出来的多值INSERT超过了 MySQL 的max_allowed_packet。解决把batchSize调小或者调大 MySQL 的max_allowed_packet比如 16M 或 32M。更稳的做法是在代码里估算单批字节数超过阈值就提前 flush而不是等数量触发。4.5 现象Checkpoint 越来越大恢复越来越慢原因是ListState里攒的数据在 Checkpoint 时被完整快照如果batchSize很大且 Checkpoint 间隔短状态会膨胀。解决控制batchSize上限Checkpoint 间隔不要短于intervalMs的两倍保证 Checkpoint 时大部分批次已经 flush 清空。另外确认状态后端用的是 RocksDB 而不是 Heap大批量 ListState 用 Heap 容易 OOM。5. 进阶用两阶段提交把“至少一次”收紧到“精确一次”前面所有方案默认是“至少一次”靠 MySQL 幂等兜底。如果你的业务不能接受任何重复比如扣款流水那就得走两阶段提交。Flink 提供了TwoPhaseCommitSinkFunction思路是数据先写到 MySQL 的临时表或带事务标记的表等 Checkpoint 完成后再提交事务失败则回滚。实现上要重写beginTransaction、preCommit、commit、abort四个方法把 JDBC 事务和 Flink 的 Checkpoint 生命周期对齐。// 两阶段提交的骨架事务在 preCommit 挂起commit 才真正生效 public class MysqlTwoPhaseSink extends TwoPhaseCommitSinkFunctionString, Connection, Void { public MysqlTwoPhaseSink() { super(new KryoSerializer(Connection.class, new ExecutionConfig()), VoidSerializer.INSTANCE); } Override protected Connection beginTransaction() throws Exception { Connection conn DriverManager.getConnection(url, user, pwd); conn.setAutoCommit(false); // 开启事务不自动提交 return conn; } Override protected void invoke(Connection conn, String value, Context context) throws Exception { // 在事务内执行写入此时数据对其他会话不可见 try (PreparedStatement ps conn.prepareStatement(insertSql)) { ps.setString(1, value); ps.executeUpdate(); } } Override protected void preCommit(Connection conn) throws Exception { // Checkpoint 前的预提交这里不 commit只做 flush conn.commit(); // 简化写法生产建议用 XA 或临时表 } Override protected void commit(Connection conn) { // Checkpoint 完成后真正提交 try { conn.commit(); conn.close(); } catch (Exception ignored) {} } Override protected void abort(Connection conn) { // 失败回滚 try { conn.rollback(); conn.close(); } catch (Exception ignored) {} } }逻辑说明beginTransaction开启一个不自动提交的连接invoke在事务内写数据preCommit在 Checkpoint 前把事务挂起commit在 Checkpoint 完成后提交abort在失败时回滚。参数说明这套方案对 MySQL 的隔离级别有要求建议用READ COMMITTED并且事务不能跨 Checkpoint 太久否则会长时间持锁。实际落地时更常见的折中是“幂等写入 唯一索引”两阶段提交只在强一致场景用因为它会显著降低吞吐。我自己的习惯是先用攒批加幂等把链路跑稳观察一周的重复率和延迟分布确认幂等真的兜住了再考虑要不要上两阶段提交。多数业务里ON DUPLICATE KEY UPDATE加业务唯一键已经够用硬上两阶段提交反而把吞吐拖下来得不偿失。希望帮到你。本文还有配套的精品资源点击获取