☰
Flink实时流上实现准批量聚合写入MySQL的生产实践
2026/10/3 2:45:23 网站建设 项目流程

简介:本资源是一套基于Flink实现Kafka实时数据流批量聚合并写入MySQL的完整工程实践方案,面向大数据开发工程师、实时计算初学者及高校相关课程学习者,解决流式数据在定时或按条数触发条件下高效聚合与关系型数据库持久化的典型问题。压缩包共9个文件,含4个核心Java代码(涵盖Flink消费、聚合、JDBC写入逻辑)、2个SQL建表与初始化脚本、1个ZooKeeper安装包(3.4.11)、1个Kafka安装包(0.9.0.0)及1个Maven配置xml,整体67.84MB,结构紧凑,覆盖环境搭建、代码开发与数据验证全链路。已有3418人学习下载,提供可直接运行的端到端示例,包含FlinkKafkaConsumer配置细节、时间/数量双触发机制实现、MySQL批量插入优化及配套依赖管理,助读者快速掌握实时数仓中关键的数据接入与落地环节。

1. Flink实时读取Kafka数据批量聚合(定时/按数量)写入MySQL:这不是“流处理入门demo”,而是生产级批流一体落地的最小可行闭环

你手头这个Flink实时读取Kafka数据批量聚合(定时按数量)写入Mysql.rar压缩包,表面看是个教学示例,实则藏着一套被反复验证过的、能直接抠出来跑通生产环境的轻量级批流融合骨架。它不依赖Flink SQL、不引入Hive Catalog、不走CDC全量+增量双链路——就用最朴素的 DataStream API + Kafka Consumer + JDBC Sink,把「每5秒或满1000条就触发一次聚合写入MySQL」这件事,从概念落到可调试、可监控、可压测的代码行里。适合正在搭建实时数仓宽表层、用户行为汇总表、IoT设备心跳统计模块的工程师;也适合刚学完Flink窗口但卡在“怎么让结果真正落库”的同学——因为这里没有玄学配置,只有三处关键参数:allowedLateness(0)、trigger(CountTrigger.of(1000))和JDBCOutputFormat.setQuery("INSERT INTO ... ON DUPLICATE KEY UPDATE ...")。压缩包里带的zookeeper-3.4.11.tar.gz和kafka_2.10-0.9.0.0.tgz虽然版本较老(对应Flink 1.7~1.9生态),但恰恰说明它避开了Flink 1.14+ 的State TTL自动清理陷阱和Kafka 3.x的SASL/SSL握手黑匣子,是那种“搭好ZK→启Kafka→起Flink集群→跑main方法→查MySQL表”四步就能看到数据进来的硬核实战包。别被标题里的“批量聚合”误导——它不是离线批处理,而是用Flink的基于事件时间的滚动窗口 + 可配置触发器,在流上模拟出可控节奏的“准批量”输出,既保实时性,又控DB写压。


2. 为什么选DataStream API而非Table/SQL?从源码结构看Flink-Kafka-Mysql链路的可控性设计

这个项目没用Flink SQL,也没用Table API的executeSql("INSERT INTO mysql_sink SELECT ..."),而是坚持用StreamExecutionEnvironment+FlinkKafkaConsumer+WindowedStream+JDBCOutputFormat的纯DataStream路径。这不是技术怀旧,而是为三个现实问题留出精准干预空间:乱序容忍粒度、窗口触发时机、JDBC写入幂等性。我们一层层拆解它的src/目录结构和核心选型逻辑。

2.1 源码包结构与各组件职责边界

压缩包解压后目录如下(已剔除IDE配置和target):

kafkasink2mysql/ ├── src/ │ ├── main/ │ │ ├── java/ │ │ │ └── com/example/flink/kafka2mysql/ │ │ │ ├── KafkaToMysqlJob.java ← 主程序入口,构建执行图 │ │ │ ├── StudentAggFunction.java ← 自定义AggregateFunction,实现sum/count逻辑 │ │ │ └── MysqlJdbcSink.java ← 封装JDBCOutputFormat,含重试+ON DUPLICATE KEY │ │ └── resources/ │ │ └── kafka.properties ← bootstrap.servers, group.id, auto.offset.reset │ ├── test/ │ └── pom.xml ├── Student.sql ← MySQL建表语句(含主键+唯一索引) ├── zookeeper-3.4.11.tar.gz ├── kafka_2.10-0.9.0.0.tgz └── README.md

提示:Student.sql中CREATE TABLE student_agg (id VARCHAR(64) PRIMARY KEY, total_score BIGINT, count INT, last_update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP)这条建表语句是关键——PRIMARY KEY保证了后续INSERT ... ON DUPLICATE KEY UPDATE能生效,last_update_time的ON UPDATE CURRENT_TIMESTAMP让你能一眼看出哪条记录是最新聚合结果。

2.2 Kafka Consumer配置:为什么用0.9.0.0版Kafka客户端?

项目pom.xml中Kafka依赖为:

<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_2.11</artifactId> <version>1.7.2</version> </dependency>

对应Kafka客户端版本为0.10.x(兼容0.9.0.0服务端)。这个组合规避了两个高频翻车点:

  • Flink 1.12+ 的Kafka 2.8+ connector默认启用partition.discovery.interval.ms=30000:在测试环境单Broker时,该参数会导致Consumer反复重平衡,日志刷屏Revoke previously assigned partitions...。而0.9.0.0版client无此机制,auto.offset.reset=earliest一设就稳。
  • Kafka 0.9.0.0的offsets.topic.replication.factor=1:避免在单节点ZooKeeper下因副本数不足导致__consumer_offsetstopic创建失败,进而使Flink任务卡在RUNNING但无数据消费。

实际启动Kafka前,必须手动修改config/server.properties:

# 必须显式设置,否则0.9.0.0默认为-1(无效值) offsets.topic.replication.factor=1 # 关闭自动创建topic(防脏数据) auto.create.topics.enable=false

2.3 窗口聚合策略:滚动窗口 + CountTrigger + EventTime的三角锚定

KafkaToMysqlJob.java中核心窗口定义如下:

DataStream<StudentEvent> kafkaStream = env .addSource(new FlinkKafkaConsumer<>("student-topic", new SimpleStringSchema(), props)) .assignTimestampsAndWatermarks( new BoundedOutOfOrdernessTimestampExtractor<String>(Time.seconds(5)) { @Override public long extractTimestamp(String element) { // 假设JSON字符串含"event_time": "2023-10-01 12:00:00" return parseEventTime(element); // 自定义解析逻辑 } } ); kafkaStream .map(new StudentEventMapFunction()) // 解析JSON → StudentEvent对象 .keyBy("studentId") // 按学生ID分组 .window(TumblingEventTimeWindows.of(Time.seconds(5))) // 5秒滚动窗口 .trigger(CountTrigger.of(1000)) // 满1000条提前触发 .aggregate(new StudentAggFunction()) // 聚合逻辑 .addSink(new MysqlJdbcSink()); // 写入MySQL

这段代码实现了双重触发保障:

  • 正常情况下,每5秒生成一个窗口,窗口内数据按studentId分组聚合;
  • 若某学生ID数据洪峰突至(如考试系统瞬间上报1000条成绩),CountTrigger.of(1000)会强制提前触发该窗口计算,避免5秒延迟堆积。

注意:CountTrigger是Trigger子类,它不替代窗口生命周期,而是在窗口活跃期内监听元素计数。allowedLateness(Time.seconds(0))已设为0,意味着完全不接受迟到数据——这对业务要求“强实时”的场景(如风控拦截)是合理取舍。


3. JDBC Sink的幂等写入:为什么不用JDBCAppendTableSink而手写MysqlJdbcSink?

Flink官方JDBCAppendTableSink只支持追加写入(INSERT INTO),但实时聚合结果需更新已有记录(如学生总分随新成绩持续累加)。若直接用INSERT,MySQL主键冲突会报错;若用REPLACE INTO,会先删后插,丢失last_update_time的ON UPDATE语义。本项目采用JDBCOutputFormat封装INSERT ... ON DUPLICATE KEY UPDATE,这是生产环境最稳妥的幂等方案。

3.1MysqlJdbcSink.java的核心实现与重试逻辑

public class MysqlJdbcSink extends RichSinkFunction<StudentAggResult> { private static final long serialVersionUID = 1L; private transient Connection connection; private transient PreparedStatement ps; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); connection = DriverManager.getConnection( "jdbc:mysql://localhost:3306/test?useSSL=false&serverTimezone=UTC", "root", "password" ); // 关键:ON DUPLICATE KEY UPDATE 保证幂等 String sql = "INSERT INTO student_agg (id, total_score, count, last_update_time) " + "VALUES (?, ?, ?, NOW()) " + "ON DUPLICATE KEY UPDATE " + "total_score = total_score + VALUES(total_score), " + "count = count + VALUES(count), " + "last_update_time = NOW()"; ps = connection.prepareStatement(sql); } @Override public void invoke(StudentAggResult value, Context context) throws Exception { ps.setString(1, value.getStudentId()); ps.setLong(2, value.getTotalScore()); ps.setInt(3, value.getCount()); ps.executeUpdate(); } @Override public void close() throws Exception { if (ps != null) ps.close(); if (connection != null) connection.close(); } }

这段代码有三个易被忽略的细节:

  • NOW()函数在VALUES()和ON DUPLICATE KEY UPDATE中各出现一次,确保无论插入还是更新,last_update_time都取当前数据库时间(非Flink TaskManager本地时间),避免时钟漂移导致的时间戳混乱;
  • total_score = total_score + VALUES(total_score)是增量更新,而非覆盖更新。假设窗口A聚合得total_score=85,窗口B聚合得total_score=92,最终数据库存的是85+92=177,符合业务对“累计总分”的语义;
  • invoke()方法未做异常捕获——这反而是正确设计。Flink的Checkpoint机制要求Sink必须是at-least-once语义,若此处try-catch吞掉SQLException,会导致数据丢失却无感知。真实部署时应配合Flink Web UI的Task Metrics → numRecordsOutPerSecond与MySQL的SHOW PROCESSLIST交叉验证写入速率。

3.2 MySQL连接池化改造:从单连接到HikariCP的平滑升级

原代码每次open()新建Connection,高并发下易触发MySQLmax_connections限制(默认151)。升级为HikariCP只需两步:

  1. pom.xml添加依赖:
<dependency> <groupId>com.zaxxer</groupId> <artifactId>HikariCP</artifactId> <version>4.0.3</version> </dependency>
  1. 修改MysqlJdbcSink.open():
private transient HikariDataSource dataSource; @Override public void open(Configuration parameters) throws Exception { HikariConfig config = new HikariConfig(); config.setJdbcUrl("jdbc:mysql://localhost:3306/test?useSSL=false&serverTimezone=UTC"); config.setUsername("root"); config.setPassword("password"); config.setMaximumPoolSize(20); // 根据Flink并行度调整 config.setMinimumIdle(5); config.setConnectionTimeout(30000); dataSource = new HikariDataSource(config); ps = dataSource.getConnection().prepareStatement(sql); }

血泪经验:maximumPoolSize建议设为Flink并行度 × 2。例如Flink Job并行度为4,则PoolSize=8;若设为20而并行度仅2,大量空闲连接会耗尽MySQL内存。


4. 避坑:Flink-Kafka-MySQL链路中五个必踩的“静默失败”点

这套方案看似简单,但在本地调试和小规模部署时,极易陷入“任务RUNNING但MySQL无数据”的黑匣子。以下是我在三套不同客户环境复现并定位的5个典型问题,按现象→原因→解决的结构给出可立即验证的排查步骤。

4.1 现象:Flink Web UI显示Source算子numRecordsInPerSecond=0,Kafka Topic确认有数据

原因:Kafka Consumer的group.id在kafka.properties中未配置,或配置为"",导致Flink使用默认随机group.id。而Kafka 0.9.0.0默认auto.offset.reset=latest,新group首次消费从最新offset开始,错过历史数据。
解决:

  • 检查src/main/resources/kafka.properties是否含group.id=test-flink-consumer;
  • 启动Flink任务前,用Kafka命令行确认数据存在:
    # 进入kafka_2.10-0.9.0.0/bin目录 ./kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic student-topic --from-beginning --max-messages 5

4.2 现象:窗口聚合结果写入MySQL,但last_update_time始终为'0000-00-00 00:00:00'

原因:MySQL服务器的SQL Mode包含NO_ZERO_DATE,且JDBC URL未显式关闭严格模式。
解决:

  • 登录MySQL执行SELECT @@sql_mode;,若返回含NO_ZERO_DATE,则修改MySQL配置文件my.cnf:
    [mysqld] sql_mode = STRICT_TRANS_TABLES,ERROR_FOR_DIVISION_BY_ZERO,NO_AUTO_CREATE_USER,NO_ENGINE_SUBSTITUTION
  • 或在JDBC URL末尾添加&zeroDateTimeBehavior=convertToNull;
  • 重启MySQL服务(仅改URL参数无效,必须重启)。

4.3 现象:Flink任务运行数小时后突然Failover,日志报java.lang.OutOfMemoryError: GC overhead limit exceeded

原因:TumblingEventTimeWindows.of(Time.seconds(5))的窗口状态未清理,Flink默认将所有窗口状态存于RocksDB,而allowedLateness=0未开启,导致窗口关闭后状态仍驻留内存。
解决:

  • 在pom.xml中添加RocksDB状态后端依赖:
    <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-statebackend-rocksdb_2.11</artifactId> <version>1.7.2</version> </dependency>
  • 在KafkaToMysqlJob.java开头添加:
    env.setStateBackend(new RocksDBStateBackend("file:///tmp/flink/checkpoints")); env.enableCheckpointing(60000); // 60秒checkpoint间隔

4.4 现象:MySQL表中出现重复studentId记录,ON DUPLICATE KEY UPDATE未生效

原因:Student.sql建表时未声明PRIMARY KEY或UNIQUE KEY,或MysqlJdbcSink中INSERT语句的VALUES()字段顺序与表结构不一致。
解决:

  • 执行DESCRIBE student_agg;确认id列为PRI(主键);
  • 对照MysqlJdbcSink.java中ps.setString(1, value.getStudentId()),确认value.getStudentId()返回的值类型为String且非null(空字符串""在MySQL中不等于NULL,但可能被当作不同主键);
  • 在invoke()方法开头加日志:System.out.println("Writing to DB: id=" + value.getStudentId());,确认传入值符合预期。

4.5 现象:Kafka Topic数据量激增,Flink任务backpressure状态变红,但CPU使用率低于30%

原因:MySQL写入成为瓶颈,JDBCOutputFormat单线程执行ps.executeUpdate(),无法利用多核。
解决:

  • 将Sink并行度显式设为2或4(需与KeyBy后的并行度匹配):
    .addSink(new MysqlJdbcSink()) .setParallelism(4); // 在addSink后链式调用
  • 同时在MysqlJdbcSink.open()中,将HikariCP的maximumPoolSize同步调至4×2=8,避免连接池成为新瓶颈。

5. 从本地调试到生产部署:三阶段验证法与MySQL写入性能压测技巧

这套方案的价值不在“能跑”,而在“能稳、能查、能扩”。我把它拆成三个递进阶段来验证——每个阶段都有明确的验收指标和失败回滚点,避免陷入“改一点、试半天、不知哪错了”的泥潭。

5.1 阶段一:单机闭环验证(15分钟内完成)

目标:确认数据从Kafka到MySQL的端到端链路畅通,且时间语义正确。
操作清单:

  1. 启动ZooKeeper:./zookeeper-3.4.11/bin/zkServer.sh start;
  2. 启动Kafka:./kafka_2.10-0.9.0.0/bin/kafka-server-start.sh ./kafka_2.10-0.9.0.0/config/server.properties;
  3. 创建Topic:./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic student-topic;
  4. 启动MySQL(确保Student.sql已执行);
  5. 编译并运行Flink Job:mvn clean package && flink run -c com.example.flink.kafka2mysql.KafkaToMysqlJob target/kafkasink2mysql-1.0.jar;
  6. 发送测试数据(模拟学生事件):
    echo '{"studentId":"S001","score":85,"event_time":"2023-10-01 12:00:01"}' | \ ./kafka-console-producer.sh --broker-list localhost:9092 --topic student-topic
  7. 验证指标:
    • Flink Web UI →Task Managers → Metrics → numRecordsInPerSecond > 0;
    • MySQL执行SELECT * FROM student_agg WHERE id='S001';,确认total_score=85,count=1,last_update_time为当前时间(误差<2秒);
    • 发送第二条数据{"studentId":"S001","score":92,...},再次查询,确认total_score=177,count=2。

注意:若第7步失败,立即执行./kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group test-flink-consumer --describe查看CURRENT-OFFSET是否前进。若未前进,说明Consumer未拉取数据,重点检查kafka.properties中的bootstrap.servers和group.id。

5.2 阶段二:时间语义压力测试(30分钟)

目标:验证TumblingEventTimeWindows在乱序数据下的正确性,以及CountTrigger的提前触发能力。
构造乱序数据脚本(gen_out_of_order.sh):

#!/bin/bash for i in {1..500}; do # 生成时间戳:前250条用t=0,后250条用t=10秒后,模拟严重乱序 if [ $i -le 250 ]; then ts="2023-10-01 12:00:00" else ts="2023-10-01 12:00:10" fi echo "{\"studentId\":\"S001\",\"score\":$((RANDOM%100)),\"event_time\":\"$ts\"}" | \ ./kafka-console-producer.sh --broker-list localhost:9092 --topic student-topic done

验证方法:

  • 启动Flink Job前,先清空MySQL表:TRUNCATE TABLE student_agg;;
  • 执行gen_out_of_order.sh;
  • 等待5秒(窗口周期),执行SELECT total_score, count FROM student_agg WHERE id='S001';;
  • 预期结果:count=500,total_score为500个随机数之和。若count<500,说明allowedLateness设置过严或Watermark生成异常。

5.3 阶段三:生产级写入压测(MySQL侧关键参数调优表)

当单机验证通过,需将Flink并行度提升至4,并发写入MySQL。此时MySQL的innodb_buffer_pool_size、max_connections等参数必须匹配。下表为针对不同Flink并行度的MySQL调优建议(基于8GB内存服务器):

Flink并行度HikariCPmaximumPoolSizeMySQLmax_connectionsMySQLinnodb_buffer_pool_size验证命令
2102002GSHOW VARIABLES LIKE 'max_connections';
4203004GSHOW ENGINE INNODB STATUS\G查看BUFFER POOL AND MEMORY
8405006GSELECT COUNT(*) FROM information_schema.PROCESSLIST;

压测执行步骤:

  1. 修改KafkaToMysqlJob.java中env.setParallelism(4);
  2. 按上表调整MySQL配置并重启;
  3. 使用sysbench模拟写入压力(避免干扰Flink):
    sysbench --db-driver=mysql --mysql-host=localhost --mysql-port=3306 \ --mysql-user=root --mysql-password=password --mysql-db=test \ oltp_write_only --tables=1 --table-size=1000000 --threads=32 run
  4. 观察Flink Web UI的latency指标(应<500ms)和MySQL的Threads_running(应<50)。

从那以后我每次上线新的Flink-Kafka-Mysql链路,都强制走一遍这三阶段:先用单条数据打穿链路,再用乱序数据校验时间语义,最后用sysbench压测DB水位。少走一次,就多一个半夜被PagerDuty叫醒的理由。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询