简介:这份资源面向大数据开发初学者与需要搭建实时数仓的工程师,聚焦Flink实时消费Kafka数据、按定时或数量条件批量聚合后写入MySQL的完整实现。包内共9个文件,以4个Java源码为核心,配合2个SQL建表脚本、pom.xml依赖配置,以及Kafka与Zookeeper的tgz、gz安装包,压缩包约67.84MB,便于直接搭建从消息队列到关系库的端到端环境。源码演示了FlinkKafkaConsumer实时摄入、时间窗口与计数窗口两种触发策略,以及通过JDBC或Table API将聚合结果持久化到MySQL的写法,同时附带Kafka集群所需的Zookeeper组件,省去单独寻找版本匹配的麻烦。已有3418人学习下载,适合对照代码理解流处理链路、快速复现实验并在此基础上改造为自身业务场景的实时聚合任务。
1. 从 Kafka 到 MySQL 的实时聚合:这套源码到底能省掉多少重复造轮子的时间
如果你正在做一个实时看板、订单统计或者设备上报汇总,大概率绕不开这条链路:Kafka 里源源不断进数据,Flink 消费后做聚合,最后落到 MySQL 给报表或后台查。听起来简单,真动手时你会发现坑全在细节里——Kafka 消费位点怎么管、聚合窗口按时间还是按条数、JDBC 批量写入怎么配、MySQL 连接器抛异常怎么排查。这套Flink实时读取Kafka数据批量聚合(定时按数量)写入Mysql.rar就是把这些细节打包成了一个能跑的工程,里面包含kafkasink2mysql源码目录、pom.xml、Student.sql建表脚本,以及zookeeper-3.4.11.tar.gz、kafka_2.10-0.9.0.0.tgz两个环境安装包。它适合两类人:一是刚接触 Flink 流处理、想找一个完整可运行 demo 把链路跑通的新手;二是已经会用 Flink 但每次写 Kafka 到 MySQL 都要重新翻文档、调参数的老手。下面我按实际拆包和复现的顺序,把这份资源讲透。
2. 拆开压缩包先看什么:工程结构与依赖版本核对
拿到一个.rar资源,别急着导入 IDE 就跑。先看清楚里面有什么、版本对不对,能省掉后面一半的报错。这套资源的目录结构不复杂,但每个文件都有它的位置意义。
2.1 目录清单与各文件职责
解压后大致是这样一个布局:
Flink实时读取Kafka数据批量聚合(定时按数量)写入Mysql/ ├── kafkasink2mysql/ # Flink 工程主目录 │ ├── src/ # Java 源码 │ └── pom.xml # Maven 依赖配置 ├── Student.sql # MySQL 建表脚本 ├── zookeeper-3.4.11.tar.gz # Zookeeper 安装包 └── kafka_2.10-0.9.0.0.tgz # Kafka 安装包kafkasink2mysql是核心,src下通常按main/java和main/resources分,Java 类里会有消费 Kafka 的 Source、聚合逻辑、写 MySQL 的 Sink 三段。pom.xml决定了 Flink、Kafka 连接器、JDBC 连接器的版本,这是最容易翻车的地方。Student.sql是目标表结构,先看它才能知道聚合结果要写成什么字段。两个.tar.gz和.tgz是环境包,说明作者默认你本地或测试机上还没有 Kafka 和 Zookeeper。
提示:
.rar在 Linux 或 macOS 上需要unrar或7z解压,Windows 用 WinRAR 即可。解压后先别改任何文件,保持原样跑通再动。
2.2 版本匹配:为什么 Kafka 0.9 和 Flink 连接器要对齐
这套资源里 Kafka 是kafka_2.10-0.9.0.0,Scala 版本 2.10,Kafka 版本 0.9.0.0。这个版本比较老,但好处是依赖少、启动快,适合本地验证链路。关键点在于:Flink 的 Kafka 连接器版本必须和 Kafka 服务端版本兼容。常见做法是,Flink 1.9 以前用flink-connector-kafka-0.9或0.10,Flink 1.9 以后统一用flink-connector-kafka并指定 Kafka 版本。
打开pom.xml,重点核对三处:
<!-- Flink 核心版本 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.11</artifactId> <version>1.9.0</version> </dependency> <!-- Kafka 连接器,注意 0.9 还是通用版 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka-0.9_2.11</artifactId> <version>1.9.0</version> </dependency> <!-- JDBC 连接器,写 MySQL 用 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-jdbc_2.11</artifactId> <version>1.9.0</version> </dependency>这里_2.11是 Scala 二进制版本,必须和 Flink 发行版一致。如果你本地 Kafka 换成了 2.x,连接器 artifactId 要改成flink-connector-kafka_2.11,否则消费时会报ClassNotFoundException或版本不兼容的NoSuchMethodError。flink-connector-jdbc在 1.9 里还不是官方一等公民,有些工程会用自定义RichSinkFunction加PreparedStatement批量提交,效果一样,但参数要自己控。
2.3 建表脚本 Student.sql 里藏着的字段约定
Student.sql不只是建个表,它定义了聚合结果的落地格式。典型内容类似:
CREATE TABLE `student_agg` ( `id` BIGINT(20) NOT NULL AUTO_INCREMENT, `class_id` VARCHAR(64) DEFAULT NULL, `stu_count` INT(11) DEFAULT '0', `window_end` DATETIME DEFAULT NULL, PRIMARY KEY (`id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;class_id是聚合维度,stu_count是聚合结果,window_end标记这批数据属于哪个时间窗口。写 Sink 时 SQL 的字段顺序、类型必须和这里一致,否则 JDBC 批量插入会报Data truncation或Column count doesn't match。常见做法是先在 MySQL 里执行这个脚本,再用DESC student_agg;确认字段类型,尤其是DATETIME和VARCHAR的长度。
3. 把链路跑起来:Kafka 生产、Flink 消费聚合、MySQL 落库
环境包和源码都齐了,接下来按数据流向一步步搭。这一章是整篇的核心,每一步都给出可抄的命令和代码片段,参数怎么改、为什么这么改一并说清。
3.1 启动 Zookeeper 与 Kafka 并造测试数据
先解压两个环境包,启动 Zookeeper 和 Kafka。Kafka 0.9 依赖 Zookeeper 存元数据,所以顺序不能反。
# 解压 tar -zxvf zookeeper-3.4.11.tar.gz tar -zxvf kafka_2.10-0.9.0.0.tgz # 启动 Zookeeper(默认 2181 端口) cd zookeeper-3.4.11 cp conf/zoo_sample.cfg conf/zoo.cfg bin/zkServer.sh start # 启动 Kafka(默认 9092 端口) cd ../kafka_2.10-0.9.0.0 bin/kafka-server-start.sh config/server.properties &启动后建一个测试 topic,并用控制台生产者往里发几条 JSON 数据,模拟学生上报:
# 建 topic,1 分区 1 副本,本地测试够用 bin/kafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 1 --topic student_topic # 开一个生产者,手动输入几条 bin/kafka-console-producer.sh --broker-list localhost:9092 --topic student_topic {"classId":"C001","stuName":"张三"} {"classId":"C001","stuName":"李四"} {"classId":"C002","stuName":"王五"}这里classId是后面聚合的 key,stuName用来计数。Kafka 0.9 的控制台生产者不支持--property parse.key=true那种键值分离,所以 key 直接放在 JSON 里,Flink 端解析后keyBy。
注意:如果 Zookeeper 启动报
JAVA_HOME相关错误,先确认echo $JAVA_HOME有值且指向 JDK 8。Kafka 0.9 对 JDK 11 支持不好,建议用 JDK 8。
3.2 Flink 消费 Kafka 的 Source 配置与反序列化
Flink 工程里消费 Kafka 的核心是FlinkKafkaConsumer09。在pom.xml依赖就绪后,Java 代码大致这样写:
// 配置 Kafka 连接参数 Properties props = new Properties(); props.setProperty("bootstrap.servers", "localhost:9092"); props.setProperty("group.id", "student_agg_group"); props.setProperty("auto.offset.reset", "earliest"); // 创建 Kafka Consumer,指定 topic 和反序列化器 FlinkKafkaConsumer09<String> kafkaSource = new FlinkKafkaConsumer09<>( "student_topic", new SimpleStringSchema(), props ); // 开启 checkpoint 时把位点提交到 Kafka,避免重复消费 kafkaSource.setCommitOffsetsOnCheckpoints(true); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); // 5 秒一次 checkpoint DataStream<String> stream = env.addSource(kafkaSource);bootstrap.servers指向 Kafka 地址,group.id是消费组,auto.offset.reset设成earliest保证第一次跑能读到历史数据。SimpleStringSchema把消息当字符串读,后面再手动解析 JSON。setCommitOffsetsOnCheckpoints(true)配合enableCheckpointing是关键:Flink 做 checkpoint 时把消费位点提交回 Kafka,任务重启后从最近一次 checkpoint 恢复,不会丢也不会重复太多。如果不开 checkpoint,位点只存在 Flink 内部,重启后行为不可控。
3.3 按时间或按数量聚合:KeyedProcessFunction 与 CountWindow 的取舍
聚合策略是这套源码的重点。摘要里提到“定时或按数量触发”,落到 Flink API 上有两种常见实现:一种是CountWindow,按条数攒够就触发;另一种是KeyedProcessFunction加定时器,按处理时间或事件时间触发。源码里大概率用的是keyBy加CountWindow或timeWindow。
按数量聚合的写法:
DataStream<Tuple2<String, Integer>> aggStream = stream .map(new MapFunction<String, Tuple2<String, Integer>>() { @Override public Tuple2<String, Integer> map(String value) throws Exception { // 解析 JSON,取出 classId,计数 1 JSONObject obj = JSON.parseObject(value); return new Tuple2<>(obj.getString("classId"), 1); } }) .keyBy(0) // 按 classId 分组 .countWindow(5) // 每 5 条触发一次 .sum(1); // 对第二个字段求和keyBy(0)按元组第一个字段分组,countWindow(5)表示每个 key 攒够 5 条就触发一次聚合,sum(1)对计数累加。这样每 5 条学生数据就会输出一个(classId, count)。如果改成定时触发,把countWindow(5)换成timeWindow(Time.minutes(1)),就是每分钟输出一次。两者可以组合,比如countWindow(5)加timeWindow的变体,但 Flink 里窗口类型不能随意叠加,常见做法是用KeyedProcessFunction自己维护计数和定时器,灵活性更高。
提示:
countWindow是滚动窗口,攒够就清空重新计数。如果你要的是“每 5 条但保留最近 10 条”这种滑动语义,得用countWindow(10, 5),参数含义是窗口大小和滑动步长。
3.4 JDBC Sink 批量写入 MySQL 与连接参数
聚合完的结果要写 MySQL。Flink 1.9 可以用flink-connector-jdbc,也可以自己写RichSinkFunction。源码里如果用的是自定义 Sink,核心逻辑是攒一批再executeBatch:
public class MysqlSink extends RichSinkFunction<Tuple2<String, Integer>> { private Connection conn; private PreparedStatement ps; private int batchSize = 100; private int count = 0; @Override public void open(Configuration parameters) throws Exception { conn = DriverManager.getConnection( "jdbc:mysql://localhost:3306/test?useSSL=false&characterEncoding=utf8", "root", "123456"); ps = conn.prepareStatement( "INSERT INTO student_agg(class_id, stu_count, window_end) VALUES(?,?,?)"); } @Override public void invoke(Tuple2<String, Integer> value, Context context) throws Exception { ps.setString(1, value.f0); ps.setInt(2, value.f1); ps.setTimestamp(3, new Timestamp(System.currentTimeMillis())); ps.addBatch(); if (++count >= batchSize) { ps.executeBatch(); conn.commit(); count = 0; } } @Override public void close() throws Exception { if (count > 0) ps.executeBatch(); ps.close(); conn.close(); } }batchSize控制多少条提交一次,设 100 到 1000 之间比较常见,太小频繁 IO,太大内存涨。useSSL=false避免 MySQL 8 以下版本 SSL 握手报错,characterEncoding=utf8防止中文乱码。close()里补一次executeBatch是血泪经验,不然最后不足一批的数据会丢。如果 MySQL 是 8.0 以上,驱动类换成com.mysql.cj.jdbc.Driver,URL 里加serverTimezone=Asia/Shanghai。
4. 避坑与排查:这套链路最容易翻车的五个地方
链路跑通一次不难,难的是稳定跑。下面五条是我在实际复现和帮人排查时遇到频率最高的,每条按现象、原因、解决写。
4.1 现象:Flink 启动就报 JDBC 连接器异常
原因:flink-connector-jdbc的版本和 Flink 核心版本不一致,或者 MySQL 驱动没打进 fat jar。常见报错是NoClassDefFoundError: com/mysql/jdbc/Driver或NoSuchMethodError。
解决:在pom.xml里显式加 MySQL 驱动依赖,版本和本地 MySQL 匹配:
<dependency> <groupId>mysql</groupId> <artifactId>mysql-connector-java</artifactId> <version>5.1.47</version> </dependency>打包时用maven-shade-plugin把依赖打进去,别用maven-assembly-plugin的默认配置,容易漏。跑之前java -jar加-cp确认驱动在 classpath 里。
4.2 现象:Kafka 消息延迟高,聚合结果半天不出来
原因:countWindow按数量触发,如果某个 key 的数据一直攒不够窗口大小,就永远不输出。比如countWindow(5),但某个班级只来了 3 条数据,这 3 条会一直卡在状态里。
解决:要么改用timeWindow保证最迟多久输出一次,要么用KeyedProcessFunction注册处理时间定时器,比如 30 秒没攒够也强制输出。源码里如果只给了countWindow,生产环境要补一个超时机制,否则冷 key 会拖垮整个作业。
4.3 现象:MySQL 里出现重复数据
原因:Flink 作业重启后从 checkpoint 恢复,但上一次 checkpoint 之后到失败前的数据被重新消费,Sink 又插了一遍。或者 Kafka 位点提交策略没配对。
解决:MySQL 表加唯一索引,比如UNIQUE KEY uk_class_window (class_id, window_end),插入用INSERT ... ON DUPLICATE KEY UPDATE。同时确认setCommitOffsetsOnCheckpoints(true)和enableCheckpointing都开了,checkpoint 间隔别设太大,5 到 10 秒比较稳。
4.4 现象:中文写入 MySQL 变成问号
原因:JDBC URL 没指定字符集,或者 MySQL 表、库的字符集是latin1。
解决:URL 加characterEncoding=utf8,建库建表用utf8mb4。Student.sql里如果写的是DEFAULT CHARSET=utf8,改成utf8mb4更保险,能存 emoji。连接后执行SHOW VARIABLES LIKE 'character%';确认。
4.5 现象:Zookeeper 或 Kafka 启动后连不上
原因:server.properties里zookeeper.connect指向的主机名解析不了,或者端口被占。Kafka 0.9 默认advertised.host.name没配,客户端拿到的是容器或内网地址。
解决:本地测试把zookeeper.connect改成localhost:2181,advertised.host.name设成localhost。用netstat -an | grep 2181和9092确认端口监听。如果之前跑过又异常退出,data目录里的myid和 Kafka 日志目录残留会导致启动失败,清掉logs和data重来。
5. 进阶:把聚合结果做成可验证、可回放的闭环
跑通一次只是开始,真正让这套资源有价值的是能验证结果对不对、能回放历史数据。我一般会加两个动作:一是用 MySQL 查询反推聚合逻辑,二是用 Kafka 重放确认幂等。
5.1 用 SQL 验证聚合结果是否符合预期
Flink 写进去的student_agg表,直接查:
SELECT class_id, SUM(stu_count) AS total, COUNT(*) AS batches FROM student_agg GROUP BY class_id ORDER BY total DESC;total应该等于 Kafka 里该班级实际发送的条数,batches是触发了几次窗口。如果对不上,先看 Kafka 里实际发了多少条,再看 Flink 日志里窗口触发次数。常见偏差是countWindow最后不足一批的数据没输出,或者 checkpoint 恢复导致重复计数。这个查询能快速定位是 Source 少读了还是 Sink 多写了。
5.2 用 Kafka 重放做幂等测试
把同一批数据再发一遍,观察 MySQL 里total是否翻倍。如果翻了,说明 Sink 没有幂等保护。解决办法是在INSERT语句里用ON DUPLICATE KEY UPDATE,配合唯一索引:
INSERT INTO student_agg(class_id, stu_count, window_end) VALUES(?,?,?) ON DUPLICATE KEY UPDATE stu_count = VALUES(stu_count);这样同一窗口重复写入只会更新计数,不会新增行。测试时把window_end固定成同一个值,重放两次,查COUNT(*)不变就说明幂等生效。
5.3 参数调优的边界:batchSize、checkpoint 间隔、窗口大小
这三个参数互相牵制。batchSize大,MySQL 压力小但延迟高;checkpoint间隔短,恢复快但开销大;窗口大,结果平滑但实时性差。我一般这样起步:batchSize=500,checkpoint=10s,countWindow=100或timeWindow=30s,然后根据 MySQL 写入 QPS 和 Flink 反压指标微调。反压看 Flink UI 的BackPressure面板,如果 Sink 是红色,先降batchSize或加 MySQL 连接池。
从那以后我每次拿到这类实时链路资源,都强制先跑一遍最小闭环:一条 Kafka 消息、一次窗口触发、一行 MySQL 记录,确认端到端通了再改参数。这套Flink实时读取Kafka数据批量聚合(定时按数量)写入Mysql.rar把环境包和源码放在一起,省掉了找版本、配依赖的时间,适合拿来当起点。希望帮到你。
本文还有配套的精品资源,点击获取