简介这份资源面向大数据实时处理方向的开发者与学习者聚焦Flink从Kafka实时消费数据、按定时或数量阈值批量聚合后写入MySQL的完整实现适合已具备Java与SQL基础、希望打通流处理链路的中级工程师参考。压缩包共9个文件约67.84MB包含4个Java源码文件、2个SQL脚本、1个XML配置以及Kafka与Zookeeper的安装包覆盖从环境搭建到业务代码的完整环节。已有3418人学习下载说明该场景在实际项目中具有较高参考价值。读者可从中获得FlinkKafkaConsumer实时摄入、定时与按数量双触发聚合策略、JDBC写入MySQL等关键代码实现并借助附带的Kafka、Zookeeper安装包快速搭建本地实时数据处理环境对照SQL脚本完成建表与结果验证从而掌握流式聚合落库的完整思路与排错方法。1. 从 Kafka 到 MySQL 的实时聚合这套源码包到底解决了什么很多做实时数仓的团队都遇到过这种尴尬Kafka 里数据哗哗地进下游 MySQL 却扛不住一条一条地写连接池被打满、主键冲突、延迟飙升最后整个链路雪崩。这套Flink实时读取Kafka数据批量聚合定时按数量写入Mysql.rar就是冲着这个痛点来的——它把 Flink 的FlinkKafkaConsumer、窗口聚合、FlinkJDBCOutputFormat串成一条完整链路并且给出了定时触发和按数量触发两种批量落库策略。压缩包里除了kafkasink2mysql源码工程和pom.xml还附带了zookeeper-3.4.11.tar.gz、kafka_2.10-0.9.0.0.tgz以及Student.sql等于把 Kafka 环境搭建、测试数据表、Flink 作业代码一次性配齐。适合正在做实时指标落库、又不想从零踩 JDBC 连接器坑的工程师尤其是还在用 Kafka 0.9 老版本、需要快速验证链路的场景。2. 环境搭建Zookeeper、Kafka 与 MySQL 的最小可用组合2.1 为什么包里锁定了 Kafka 0.9 和 Zookeeper 3.4.11压缩包里的kafka_2.10-0.9.0.0.tgz和zookeeper-3.4.11.tar.gz不是随便塞的。Kafka 0.9 是最后一个强依赖 Zookeeper 做消费者位移存储的大版本之后 0.10 开始逐步把位移迁到 Kafka 内部 topic。这套源码用的FlinkKafkaConsumer在对接 0.9 时需要 Zookeeper 参与 offset 管理所以两个包必须版本对齐。如果你手头是 Kafka 2.x 或 3.x直接换包会报NoSuchMethodError或ClassNotFoundException因为FlinkKafkaConsumer09和FlinkKafkaConsumer的构造参数、KafkaDeserializationSchema接口都不一样。常见做法是先按包内版本把环境跑通再根据生产集群版本替换pom.xml里的flink-connector-kafka-0.9_2.10依赖同时把消费者类名从FlinkKafkaConsumer09改成对应版本。2.2 三步把 Kafka 和 MySQL 拉起来先解压 Zookeeper 和 Kafka启动顺序不能反Zookeeper 没起来 Kafka 会直接退出。# 启动 Zookeeper默认 2181 端口 tar -zxvf zookeeper-3.4.11.tar.gz cd zookeeper-3.4.11/conf cp zoo_sample.cfg zoo.cfg cd ../bin ./zkServer.sh start # 启动 Kafka依赖上面的 Zookeeper tar -zxvf kafka_2.10-0.9.0.0.tgz cd kafka_2.10-0.9.0.0 # 后台启动指定 Zookeeper 地址 bin/kafka-server-start.sh -daemon config/server.properties # 建一个测试 topic分区数给 1 方便观察 bin/kafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 1 --topic student_topic逻辑说明zkServer.sh start拉起协调服务Kafka 的server.properties里zookeeper.connectlocalhost:2181默认指向本机。建 topic 时--partitions 1是为了让 Flink 的并行度与分区数匹配避免某些 subtask 空转。参数上--replication-factor 1单机测试够用生产至少 2。MySQL 这边直接导入包里的Student.sql表结构通常是id、name、age这类字段用来承接聚合结果。# 登录 MySQL 后执行 source /path/to/Student.sql; # 确认表已建好 show tables;注意Student.sql里如果用了ENGINEInnoDB和utf8字符集导入前先确认 MySQL 版本支持5.7 和 8.0 在默认字符集上有差异8.0 建议改成utf8mb4。2.3 源码工程kafkasink2mysql的依赖梳理pom.xml是整条链路的命门。核心依赖一般包括flink-streaming-java、flink-connector-kafka-0.9、flink-jdbc以及 MySQL 驱动。版本上Flink 1.9 之前flink-jdbc是独立包1.9 之后推荐用flink-connector-jdbc。这套源码大概率基于 Flink 1.9 或更早因为FlinkJDBCOutputFormat在 1.9 被标记废弃1.11 之后彻底移除。如果你用新版 Flink需要把 sink 部分改成JdbcSink.sink()否则编译直接报cannot resolve symbol。!-- pom.xml 关键片段版本按实际集群调整 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka-0.9_2.10/artifactId version1.9.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-jdbc_2.10/artifactId version1.9.0/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version5.1.47/version /dependency参数说明flink-connector-kafka-0.9_2.10里的2.10是 Scala 版本必须和集群里 Flink 的 Scala 版本一致否则运行时报NoClassDefFoundError。MySQL 驱动 5.1.47 对应 MySQL 5.x连 8.0 要换8.0.28并改驱动类名为com.mysql.cj.jdbc.Driver。3. 核心逻辑定时与按数量两种批量聚合怎么落地3.1 从 Kafka 读取数据的反序列化与水位线Flink 读 Kafka 的第一步是定义DeserializationSchema把字节数组转成业务对象。源码里通常用SimpleStringSchema或自定义KafkaDeserializationSchema。如果数据是 JSON自定义 schema 里用ObjectMapper解析解析失败要返回null而不是抛异常否则整个作业会挂。// 自定义反序列化把 Kafka 消息转成 Student 对象 public class StudentDeserializationSchema implements DeserializationSchemaStudent { private ObjectMapper mapper new ObjectMapper(); Override public Student deserialize(byte[] message) { try { return mapper.readValue(message, Student.class); } catch (IOException e) { // 脏数据直接丢弃避免作业失败 return null; } } Override public boolean isEndOfStream(Student nextElement) { return false; } Override public TypeInformationStudent getProducedType() { return TypeInformation.of(Student.class); } }逻辑说明deserialize里捕获异常返回nullFlink 会在后续算子中过滤掉。isEndOfStream对流式数据永远返回false。getProducedType告诉 Flink 类型信息避免 Kryo 序列化带来的性能损耗。参数上ObjectMapper建议注册JavaTimeModule处理时间字段否则遇到日期格式会抛InvalidDefinitionException。水位线Watermark的设置决定了聚合窗口何时触发。如果数据带事件时间用BoundedOutOfOrdernessTimestampExtractor最大乱序时间设 5 秒到 30 秒之间具体看业务延迟。DataStreamStudent stream env.addSource( new FlinkKafkaConsumer09(student_topic, new StudentDeserializationSchema(), properties)) .assignTimestampsAndWatermarks( new BoundedOutOfOrdernessTimestampExtractorStudent(Time.seconds(10)) { Override public long extractTimestamp(Student element) { return element.getEventTime(); } });参数说明Time.seconds(10)是允许的最大乱序设太小会丢迟到数据设太大会增加窗口延迟。extractTimestamp返回毫秒级时间戳如果数据里是秒级要乘 1000。3.2 定时触发用窗口实现每 5 分钟落一次库定时聚合靠TumblingProcessingTimeWindows或TumblingEventTimeWindows。处理时间窗口简单不依赖数据里的时间戳事件时间窗口准确但要处理迟到数据。源码里大概率用处理时间因为测试方便。// 每 5 分钟一个窗口对 age 求和 DataStreamStudentAgg aggStream stream .keyBy(Student::getName) .window(TumblingProcessingTimeWindows.of(Time.minutes(5))) .aggregate(new AggregateFunctionStudent, StudentAgg, StudentAgg() { Override public StudentAgg createAccumulator() { return new StudentAgg(); } Override public StudentAgg add(Student value, StudentAgg acc) { acc.setName(value.getName()); acc.setTotalAge(acc.getTotalAge() value.getAge()); acc.setCount(acc.getCount() 1); return acc; } Override public StudentAgg getResult(StudentAgg acc) { return acc; } Override public StudentAgg merge(StudentAgg a, StudentAgg b) { a.setTotalAge(a.getTotalAge() b.getTotalAge()); a.setCount(a.getCount() b.getCount()); return a; } });逻辑说明keyBy按姓名分组同一个人数据进同一个窗口。aggregate是增量聚合比apply全量缓存更省内存。createAccumulator初始化累加器add逐条累加getResult输出最终结果merge在会话窗口或合并窗口时用滚动窗口其实用不到但必须实现。参数上Time.minutes(5)改成Time.seconds(30)就是半分钟一次按业务 QPS 调整。3.3 按数量触发用CountTrigger控制每 1000 条写一次按数量触发需要CountTrigger配合WindowAssigner。但注意CountTrigger在并行度大于 1 时是每个 subtask 独立计数不是全局计数。如果要求全局精确 1000 条得先keyBy到一个并行度为 1 的算子再计数或者用GlobalWindow加自定义触发器。// 每 1000 条触发一次窗口计算 DataStreamStudentAgg countStream stream .keyBy(Student::getName) .window(TumblingProcessingTimeWindows.of(Time.minutes(5))) .trigger(CountTrigger.of(1000)) .aggregate(new AggregateFunctionStudent, StudentAgg, StudentAgg() { // 同上省略 });逻辑说明CountTrigger.of(1000)表示窗口内元素达到 1000 条就触发一次计算同时窗口结束时间到了也会触发。这样兼顾了数量和定时。参数上1000 这个值要根据 MySQL 的max_allowed_packet和单次插入行数调整一般 500 到 2000 之间比较稳。如果设成 10000单次INSERT语句可能超过 1MBMySQL 直接拒绝。3.4 写入 MySQLFlinkJDBCOutputFormat的配置与批量提交FlinkJDBCOutputFormat是这套源码的落库核心。它内部用PreparedStatement.addBatch()和executeBatch()实现批量提交比一条一条executeUpdate快一个数量级。// 构建 JDBC OutputFormat FlinkJDBCOutputFormat outputFormat FlinkJDBCOutputFormat.buildJDBCOutputFormat() .setDrivername(com.mysql.jdbc.Driver) .setDBUrl(jdbc:mysql://localhost:3306/test?useSSLfalse) .setUsername(root) .setPassword(123456) .setQuery(INSERT INTO student_agg(name, total_age, cnt) VALUES(?,?,?)) .setSqlTypes(new int[]{Types.VARCHAR, Types.INTEGER, Types.INTEGER}) .setBatchInterval(1000) // 每 1000 条提交一次 .finish(); aggStream.addSink(outputFormat);逻辑说明setQuery里的占位符顺序必须和setSqlTypes一一对应否则报SQLException: Parameter index out of range。setBatchInterval(1000)是 JDBC 层的批量提交阈值和窗口的CountTrigger是两回事——窗口触发后可能一次来 1000 条JDBC 再攒 1000 条提交实际落库可能是 2000 条一次。参数上useSSLfalse在测试环境省去证书配置生产建议开启。setSqlTypes用java.sql.Types常量VARCHAR对应StringINTEGER对应int。4. 避坑排查这套链路最容易翻车的五个地方4.1 现象作业启动后 Kafka 消费不到数据日志显示Offset commit failed原因Kafka 0.9 的 offset 默认存在 Zookeeper但FlinkKafkaConsumer09的enable.auto.commit如果设为true会和 Flink 的 checkpoint 机制冲突导致 offset 提交混乱。解决把properties.setProperty(enable.auto.commit, false)让 Flink 在 checkpoint 完成时统一提交 offset。同时确认zookeeper.connect和 Kafka 的zookeeper.connect一致。4.2 现象MySQL 报Packet for query is too large原因setBatchInterval设得太大或者窗口内聚合结果太多单次executeBatch的数据量超过 MySQL 的max_allowed_packet默认 4MB。解决把setBatchInterval降到 500 以下或者调大 MySQL 的max_allowed_packet16M。更稳妥的做法是在 JDBC URL 里加rewriteBatchedStatementstrue让驱动把多条INSERT合并成一条减少网络往返。4.3 现象Flink 作业频繁重启日志里java.net.SocketTimeoutException原因MySQL 连接被防火墙或wait_timeout断开FlinkJDBCOutputFormat持有的连接失效。解决在 JDBC URL 里加autoReconnecttruefailOverReadOnlyfalse同时把 MySQL 的wait_timeout调到 8 小时以上。如果用的是连接池确保validationQuery设为SELECT 1。4.4 现象聚合结果重复写入MySQL 里出现多条相同记录原因Flink 的 at-least-once 语义下checkpoint 失败恢复后窗口会重新计算导致重复落库。解决把 MySQL 表的主键设为name加窗口时间插入语句改成INSERT ... ON DUPLICATE KEY UPDATE用幂等写入抵消重复。或者开启 Flink 的 exactly-once但这需要 Kafka 和 MySQL 都支持两阶段提交配置复杂度高。4.5 现象FlinkJDBCOutputFormat编译报错cannot resolve symbol原因Flink 1.11 之后flink-jdbc被移除FlinkJDBCOutputFormat类不存在。解决改用flink-connector-jdbc的JdbcSink.sink()或者把 Flink 版本降到 1.9。如果必须用新版sink 代码要重写// Flink 1.11 的 JDBC Sink 写法 JdbcSink.sink( INSERT INTO student_agg(name, total_age, cnt) VALUES(?,?,?), (ps, studentAgg) - { ps.setString(1, studentAgg.getName()); ps.setInt(2, studentAgg.getTotalAge()); ps.setInt(3, studentAgg.getCount()); }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchIntervalMs(5000) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(jdbc:mysql://localhost:3306/test) .withUsername(root) .withPassword(123456) .build() );参数说明withBatchSize(1000)对应条数触发withBatchIntervalMs(5000)对应定时触发两者满足其一就提交。这正好复现了源码里定时加按数量的双策略。5. 进阶技巧用 Flink 火焰图定位聚合算子的性能瓶颈链路跑通之后下一步是调优。我一般会先看 Flink 的火焰图确认时间花在反序列化、聚合还是 JDBC 写入上。火焰图入口在 Flink Web UI 的 Job 页面点击算子后的Flame Graph链接或者直接访问http://jobmanager:8081/#/job/{jobId}/flamegraph。如果反序列化占了大头说明ObjectMapper创建太频繁把它提成静态成员或者用RichFlatMapFunction的open方法初始化。如果 JDBC 写入占了大头先看setBatchInterval是不是太小再确认 MySQL 的innodb_flush_log_at_trx_commit是不是 1改成 2 能提升写入吞吐但会牺牲一点持久性。另一个技巧是给 Kafka 消费者加partitionDiscoveryIntervalMillis默认 5 分钟如果 topic 分区动态增加设成 30 秒能更快感知。还有FlinkKafkaConsumer09的setStartFromGroupOffsets和setStartFromEarliest的区别前者从 Zookeeper 记录的 offset 开始后者从最早数据开始。测试阶段用setStartFromEarliest方便复现生产必须用setStartFromGroupOffsets否则每次重启都全量消费。验证聚合结果是否正确除了查 MySQL还可以在 Flink Web UI 的Accumulators里看numRecordsIn和numRecordsOut如果numRecordsOut远小于numRecordsIn说明聚合生效了。我习惯在AggregateFunction的getResult里加一行LOG.info把每个窗口的结果打出来和 MySQL 里的记录对一遍确认没有漏写或重复。从那以后我每次搭 Flink 到 MySQL 的链路都强制先跑一遍Student.sql建表、再确认pom.xml里 Flink 版本和连接器版本对齐、最后用火焰图看一遍算子耗时这三步走完基本不会翻车。希望帮到你。本文还有配套的精品资源点击获取