资讯动态

Flink实时读取Kafka数据批量聚合写入MySQL实战源码解析

发布时间:2026/10/5 4:52:42 来源:尧图企业网站定制
简介这份资源面向大数据开发初学者与需要搭建实时数仓的工程师聚焦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.0Scala 版本 2.10Kafka 版本 0.9.0.0。这个版本比较老但好处是依赖少、启动快适合本地验证链路。关键点在于Flink 的 Kafka 连接器版本必须和 Kafka 服务端版本兼容。常见做法是Flink 1.9 以前用flink-connector-kafka-0.9或0.10Flink 1.9 以后统一用flink-connector-kafka并指定 Kafka 版本。打开pom.xml重点核对三处!-- Flink 核心版本 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.11/artifactId version1.9.0/version /dependency !-- Kafka 连接器注意 0.9 还是通用版 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka-0.9_2.11/artifactId version1.9.0/version /dependency !-- JDBC 连接器写 MySQL 用 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc_2.11/artifactId version1.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) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;class_id是聚合维度stu_count是聚合结果window_end标记这批数据属于哪个时间窗口。写 Sink 时 SQL 的字段顺序、类型必须和这里一致否则 JDBC 批量插入会报Data truncation或Column count doesnt 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 数据模拟学生上报# 建 topic1 分区 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是后面聚合的 keystuName用来计数。Kafka 0.9 的控制台生产者不支持--property parse.keytrue那种键值分离所以 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 和反序列化器 FlinkKafkaConsumer09String kafkaSource new FlinkKafkaConsumer09( student_topic, new SimpleStringSchema(), props ); // 开启 checkpoint 时把位点提交到 Kafka避免重复消费 kafkaSource.setCommitOffsetsOnCheckpoints(true); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); // 5 秒一次 checkpoint DataStreamString 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。按数量聚合的写法DataStreamTuple2String, Integer aggStream stream .map(new MapFunctionString, Tuple2String, Integer() { Override public Tuple2String, 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核心逻辑是攒一批再executeBatchpublic class MysqlSink extends RichSinkFunctionTuple2String, 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?useSSLfalsecharacterEncodingutf8, root, 123456); ps conn.prepareStatement( INSERT INTO student_agg(class_id, stu_count, window_end) VALUES(?,?,?)); } Override public void invoke(Tuple2String, 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太大内存涨。useSSLfalse避免 MySQL 8 以下版本 SSL 握手报错characterEncodingutf8防止中文乱码。close()里补一次executeBatch是血泪经验不然最后不足一批的数据会丢。如果 MySQL 是 8.0 以上驱动类换成com.mysql.cj.jdbc.DriverURL 里加serverTimezoneAsia/Shanghai。4. 避坑与排查这套链路最容易翻车的五个地方链路跑通一次不难难的是稳定跑。下面五条是我在实际复现和帮人排查时遇到频率最高的每条按现象、原因、解决写。4.1 现象Flink 启动就报 JDBC 连接器异常原因flink-connector-jdbc的版本和 Flink 核心版本不一致或者 MySQL 驱动没打进 fat jar。常见报错是NoClassDefFoundError: com/mysql/jdbc/Driver或NoSuchMethodError。解决在pom.xml里显式加 MySQL 驱动依赖版本和本地 MySQL 匹配dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version5.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 加characterEncodingutf8建库建表用utf8mb4。Student.sql里如果写的是DEFAULT CHARSETutf8改成utf8mb4更保险能存 emoji。连接后执行SHOW VARIABLES LIKE character%;确认。4.5 现象Zookeeper 或 Kafka 启动后连不上原因server.properties里zookeeper.connect指向的主机名解析不了或者端口被占。Kafka 0.9 默认advertised.host.name没配客户端拿到的是容器或内网地址。解决本地测试把zookeeper.connect改成localhost:2181advertised.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间隔短恢复快但开销大窗口大结果平滑但实时性差。我一般这样起步batchSize500checkpoint10scountWindow100或timeWindow30s然后根据 MySQL 写入 QPS 和 Flink 反压指标微调。反压看 Flink UI 的BackPressure面板如果 Sink 是红色先降batchSize或加 MySQL 连接池。从那以后我每次拿到这类实时链路资源都强制先跑一遍最小闭环一条 Kafka 消息、一次窗口触发、一行 MySQL 记录确认端到端通了再改参数。这套Flink实时读取Kafka数据批量聚合定时按数量写入Mysql.rar把环境包和源码放在一起省掉了找版本、配依赖的时间适合拿来当起点。希望帮到你。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑