
简介这份资源面向大数据实时处理方向的开发者与学习者聚焦Flink从Kafka实时读取数据、按定时或数量条件批量聚合后写入MySQL的完整实现适合已具备Java与SQL基础、希望打通流处理链路的中级读者参考。压缩包共9个文件以4个Java源码为核心配合2个SQL建表脚本、1个pom.xml依赖配置以及Kafka与Zookeeper的tgz、gz安装包整体约67.84MB便于一次性搭建可运行的实时数据处理环境。源码展示了FlinkKafkaConsumer实时摄入、定时与按数量两种触发策略的聚合逻辑以及通过JDBC或Table API将结果持久化到MySQL的写法读者可据此理解状态管理、容错与批量写入的配合方式。目前已有3418人学习下载适合作为实时数仓入门与项目练手的参考案例。1. Flink 按数量攒批写 MySQL一条被低估的实时链路Kafka 里每秒几万条消息下游 MySQL 却扛不住一条一条 insert——这不是理论问题是很多人第一次把 Flink 实时链路接到关系库时必然撞上的墙。标题里的方案说白了就一件事Flink 消费 Kafka不急着写库先在内存里按数量攒一批攒够阈值再一次性批量写入 MySQL。它解决的是「实时性要求没那么极致、但写入吞吐必须压下来」的场景比如订单流水落库、设备上报归档、日志明细入仓前的缓冲层。适合谁手里已经有 Kafka 和 MySQL、想用 Flink 做轻量实时同步、又不想引入 ClickHouse 这类额外组件的团队。核心关键词就三个Flink、Kafka、批量聚合外加一个容易被忽略的「定时」兜底。2. 为什么是「按数量攒批」而不是来一条写一条2.1 逐条写入 MySQL 到底慢在哪很多人第一反应是 Flink 里接个 JDBC Sink来一条写一条代码十行搞定。跑起来才发现 MySQL 的 QPS 上不去Flink 的背压一路顶回 Kafka消费延迟肉眼可见地涨。根因不在 Flink在 MySQL 的写入模型每条 insert 都是一次独立事务要写 redo log、刷 binlog、等磁盘 fsync。单条写入的瓶颈从来不是 SQL 解析而是事务提交的网络往返和日志刷盘。你把 1000 条拆成 1000 次提交和攒成 1 次提交磁盘 IO 次数差三个数量级。批量写入的本质是把 N 次事务合并成 1 次。MySQL 的INSERT INTO ... VALUES (...),(...),(...)多值语法配合 JDBC 的rewriteBatchedStatementstrue能把多条 insert 在驱动层重写成一条多值语句网络往返从 N 次降到 1 次。这是整个方案性能提升的主要来源不是 Flink 的功劳是 JDBC 批处理的功劳。2.2 攒批的两种触发条件数量和定时只按数量攒有个致命问题流量低谷时比如半夜只有零星几条数据批次永远攒不够阈值数据就卡在内存里出不去。所以标题里特意点了「定时」——数量阈值和定时器必须同时存在谁先到谁触发。常见做法是 Flink 的ProcessingTimeService注册一个周期性定时器或者用KeyedProcessFunction的onTimer每隔 N 秒强制 flush 一次当前缓冲区。这里有个选型细节按数量攒批天然适合用 Flink 的RichSinkFunction自己维护一个List在invoke里累加到阈值就批量提交。但要注意RichSinkFunction是单并发的状态如果你开了多个并行度每个 subtask 各攒各的批MySQL 那边看到的是多个并发批量写这通常没问题但要确认连接池够用。2.3 和窗口聚合的区别别用错工具有人会问Flink 不是有window吗为什么不用滚动窗口攒批因为窗口是面向「聚合计算」的它输出的是聚合结果不是原始明细。你要的是把 1000 条原始消息原样写进 MySQL窗口反而会把它们算成一个值。除非你的需求本身就是「每分钟统计一次写入」那才用窗口。这个方案要的是缓冲不是聚合别混。提示如果你的下游是 ClickHouse 而不是 MySQL批量写入的收益更明显因为 ClickHouse 对单条 insert 极不友好但对批量写入吞吐极高。热词里「flink 实现 mysql 同步到 clickhouse」本质是同一套攒批思路换了个 Sink。3. 从 Kafka 到 MySQL 的最小可跑通链路3.1 环境与依赖三个组件的最小版本组合不追求最新追求能跑通。Flink 用 1.17 或 1.18 都行Kafka 用 3.xMySQL 5.7 或 8.0 均可。依赖上Flink 侧需要flink-connector-kafka和flink-connector-jdbcMySQL 侧需要mysql-connector-java驱动。注意 Flink 1.17 之后 JDBC connector 的坐标变了别抄老教程的flink-jdbc那个已经废弃。!-- pom.xml 关键依赖版本按你的 Flink 大版本对齐 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version3.1.0-1.18/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version3.1.2-1.18/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency版本号要对齐 Flink 大版本3.1.0-1.18里的1.18就是 Flink 版本。驱动用 8.0.33 是因为它同时兼容 MySQL 5.7 和 8.0省得你为不同环境换驱动。如果公司还在用 MySQL 5.7.44这个驱动也能连但记得 URL 里关掉 SSL否则会撞上「mysql ssl 连接错误」那个经典坑。3.2 Kafka Source 的消费配置KafkaSourceString source KafkaSource.Stringbuilder() .setBootstrapServers(kafka-broker:9092) .setTopics(order-topic) .setGroupId(flink-batch-writer) .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.LATEST)) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStreamString stream env.fromSource( source, WatermarkStrategy.noWatermarks(), kafka-source );setGroupId决定消费位点存在哪换 group 会从头或从最新重来调试时注意。OffsetsInitializer.committedOffsets表示优先用已提交位点没有才按LATEST走这是生产环境的稳妥选择。noWatermarks是因为这个方案不依赖事件时间纯处理时间攒批省掉水位线的复杂度。如果 Kafka 消息延迟高先查是不是消费组里有人卡住了而不是急着调 Flink 参数。3.3 自定义攒批 Sink 的核心逻辑public class BatchMysqlSink extends RichSinkFunctionString { private transient ListString buffer; private transient Connection conn; private transient PreparedStatement ps; private final int batchSize; // 数量阈值 private final long flushInterval; // 定时阈值毫秒 public BatchMysqlSink(int batchSize, long flushInterval) { this.batchSize batchSize; this.flushInterval flushInterval; } Override public void open(Configuration params) throws Exception { buffer new ArrayList(batchSize); conn DriverManager.getConnection( jdbc:mysql://mysql-host:3306/db?rewriteBatchedStatementstrueuseSSLfalse, user, pass); conn.setAutoCommit(false); // 手动控制事务 ps conn.prepareStatement(INSERT INTO order_detail(order_id, amount) VALUES (?, ?)); // 注册定时 flush兜底低峰期数据 ScheduledExecutorService scheduler Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(this::flushSafely, flushInterval, flushInterval, TimeUnit.MILLISECONDS); } Override public void invoke(String value, Context context) throws Exception { buffer.add(value); if (buffer.size() batchSize) { flush(); } } private synchronized void flush() throws Exception { if (buffer.isEmpty()) return; for (String line : buffer) { String[] f line.split(,); ps.setString(1, f[0]); ps.setBigDecimal(2, new BigDecimal(f[1])); ps.addBatch(); } ps.executeBatch(); conn.commit(); buffer.clear(); } private void flushSafely() { try { flush(); } catch (Exception e) { /* 记录日志别吞异常 */ } } }rewriteBatchedStatementstrue是整段代码里最值钱的参数没有它addBatch只是攒在驱动内存里执行时还是一条条发。setAutoCommit(false)配合手动commit保证一批是一个事务。定时器用scheduleAtFixedRate周期触发注意flush加了synchronized因为定时线程和invoke线程会并发访问 buffer不加锁就是血泪翻车现场。flushSafely里别把异常吞掉至少打日志否则数据丢了都不知道。3.4 参数怎么定batchSize 和 flushInterval 的取值参数建议起点调整方向影响batchSize500写入慢就加大到 1000~2000越大吞吐越高内存占用和单批失败重试成本也越高flushInterval2000ms实时性要求高就降到 500ms越小延迟越低但空批提交次数变多连接池大小并行度 × 2按 subtask 数量配不够会连接超时rewriteBatchedStatementstrue别关关了批量等于白做batchSize 不是越大越好。单批 5000 条时如果一条数据格式有问题导致整批失败你要重试 5000 条。我一般从 500 起步观察 MySQL 的写入延迟和 Flink 的 checkpoint 大小再调。flushInterval 要和业务容忍的延迟对齐2 秒是个不激进也不保守的默认值。4. 避坑与排查那些让链路半夜挂掉的问题4.1 数据攒在内存里checkpoint 越来越大现象Flink 的 checkpoint 体积持续增长恢复时间越来越长。原因buffer 是 Sink 的普通成员变量没纳入 Flink 的状态管理checkpoint 时它不会被快照但 Flink 也不知道里面有多少数据恢复后这部分数据直接丢。解决要么接受「攒批数据不保证 exactly-once」这个前提要么把 buffer 做成ListState在snapshotState里持久化。多数业务能接受 at-least-once但你要清楚这个取舍。4.2 定时 flush 和 invoke 抢锁导致吞吐骤降现象加了定时器后吞吐不升反降。原因flush方法加了synchronized定时线程频繁抢锁invoke被阻塞。解决把 flushInterval 调大或者用双缓冲——一个 buffer 在写另一个在收切换时加锁。简单场景下把 flushInterval 设成 2 秒以上锁竞争基本可忽略。4.3 MySQL 连接被服务端掐断现象跑几小时后报Communications link failure。原因MySQL 的wait_timeout默认 8 小时空闲连接被服务端关掉而客户端不知道。解决URL 里加autoReconnecttrue或者用连接池HikariCP并配置maxLifetime小于 MySQL 的wait_timeout。别用DriverManager裸连跑生产这是常识但总有人踩。4.4 批量 insert 撞上主键冲突整批回滚现象一批 500 条里有一条主键重复整批executeBatch失败。原因MySQL 默认整批是一个事务一条失败全批回滚。解决SQL 改成INSERT IGNORE或ON DUPLICATE KEY UPDATE按业务语义选。如果必须严格报错那就把 batchSize 调小降低单批失败的影响面。4.5 并行度大于 1 时数据顺序错乱现象同一订单的两条消息写进 MySQL 后顺序反了。原因Kafka 只保证分区内有序Flink 多并行度消费多个分区Sink 端各写各的。解决如果业务要求同一 key 有序在 Kafka 侧按 key 分区Flink 侧用keyBy后再攒批保证同 key 进同一 subtask。不要求顺序就别折腾多数明细表不在乎。5. 进阶让攒批链路更稳的两个技巧第一个技巧是给 Sink 加一层「失败降级」。批量提交失败时不要整批丢弃而是把这一批拆成单条逐条重试把失败的那条单独记到死信表。这样既保住了吞吐又不会因为一条脏数据丢掉整批。实现上就是在flush的 catch 块里遍历 buffer逐条executeUpdate成功的跳过失败的写日志表。第二个技巧是用 Flink 的 Metrics 暴露 buffer 当前大小和最近一次 flush 耗时。buffer 长期接近 batchSize 说明下游写入跟不上flush 耗时突增说明 MySQL 有锁等待或磁盘压力。这两个指标比看 Flink UI 的背压更直接。我一般用getRuntimeContext().getMetricGroup().gauge()注册一个 gauge接到 Prometheus 上配个告警阈值。验证这套链路是否真的生效别只看「数据写进去了」。压测时对比两个数字逐条写入的 TPS 和攒批写入的 TPS正常情况后者是前者的 10 到 50 倍。如果差距不到 5 倍八成是rewriteBatchedStatements没生效或者 batchSize 设得太小。我自己踩过的最大坑是忘了在 URL 里加这个参数白调了一下午 batchSize后来看 MySQL 的 general log 才发现每条 insert 还是单独发的。希望帮到你。本文还有配套的精品资源点击获取