ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink实时流上实现准批量聚合写入MySQL的生产实践

Flink实时流上实现准批量聚合写入MySQL的生产实践 简介本资源是一套基于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 ...)而是坚持用StreamExecutionEnvironmentFlinkKafkaConsumerWindowedStreamJDBCOutputFormat的纯DataStream路径。这不是技术怀旧而是为三个现实问题留出精准干预空间乱序容忍粒度、窗口触发时机、JDBC写入幂等性。我们一层层拆解它的src/目录结构和核心选型逻辑。2.1 源码包结构与各组件职责边界压缩包解压后目录如下已剔除IDE配置和targetkafkasink2mysql/ ├── 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 groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka_2.11/artifactId version1.7.2/version /dependency对应Kafka客户端版本为0.10.x兼容0.9.0.0服务端。这个组合规避了两个高频翻车点Flink 1.12 的Kafka 2.8 connector默认启用partition.discovery.interval.ms30000在测试环境单Broker时该参数会导致Consumer反复重平衡日志刷屏Revoke previously assigned partitions...。而0.9.0.0版client无此机制auto.offset.resetearliest一设就稳。Kafka 0.9.0.0的offsets.topic.replication.factor1避免在单节点ZooKeeper下因副本数不足导致__consumer_offsetstopic创建失败进而使Flink任务卡在RUNNING但无数据消费。实际启动Kafka前必须手动修改config/server.properties# 必须显式设置否则0.9.0.0默认为-1无效值 offsets.topic.replication.factor1 # 关闭自动创建topic防脏数据 auto.create.topics.enablefalse2.3 窗口聚合策略滚动窗口 CountTrigger EventTime的三角锚定KafkaToMysqlJob.java中核心窗口定义如下DataStreamStudentEvent kafkaStream env .addSource(new FlinkKafkaConsumer(student-topic, new SimpleStringSchema(), props)) .assignTimestampsAndWatermarks( new BoundedOutOfOrdernessTimestampExtractorString(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而手写MysqlJdbcSinkFlink官方JDBCAppendTableSink只支持追加写入INSERT INTO但实时聚合结果需更新已有记录如学生总分随新成绩持续累加。若直接用INSERTMySQL主键冲突会报错若用REPLACE INTO会先删后插丢失last_update_time的ON UPDATE语义。本项目采用JDBCOutputFormat封装INSERT ... ON DUPLICATE KEY UPDATE这是生产环境最稳妥的幂等方案。3.1MysqlJdbcSink.java的核心实现与重试逻辑public class MysqlJdbcSink extends RichSinkFunctionStudentAggResult { 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?useSSLfalseserverTimezoneUTC, 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_score85窗口B聚合得total_score92最终数据库存的是8592177符合业务对“累计总分”的语义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只需两步pom.xml添加依赖dependency groupIdcom.zaxxer/groupId artifactIdHikariCP/artifactId version4.0.3/version /dependency修改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?useSSLfalseserverTimezoneUTC); 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则PoolSize8若设为20而并行度仅2大量空闲连接会耗尽MySQL内存。4. 避坑Flink-Kafka-MySQL链路中五个必踩的“静默失败”点这套方案看似简单但在本地调试和小规模部署时极易陷入“任务RUNNING但MySQL无数据”的黑匣子。以下是我在三套不同客户环境复现并定位的5个典型问题按现象→原因→解决的结构给出可立即验证的排查步骤。4.1 现象Flink Web UI显示Source算子numRecordsInPerSecond0Kafka Topic确认有数据原因Kafka Consumer的group.id在kafka.properties中未配置或配置为导致Flink使用默认随机group.id。而Kafka 0.9.0.0默认auto.offset.resetlatest新group首次消费从最新offset开始错过历史数据。解决检查src/main/resources/kafka.properties是否含group.idtest-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 54.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末尾添加zeroDateTimeBehaviorconvertToNull重启MySQL服务仅改URL参数无效必须重启。4.3 现象Flink任务运行数小时后突然Failover日志报java.lang.OutOfMemoryError: GC overhead limit exceeded原因TumblingEventTimeWindows.of(Time.seconds(5))的窗口状态未清理Flink默认将所有窗口状态存于RocksDB而allowedLateness0未开启导致窗口关闭后状态仍驻留内存。解决在pom.xml中添加RocksDB状态后端依赖dependency groupIdorg.apache.flink/groupId artifactIdflink-statebackend-rocksdb_2.11/artifactId version1.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×28避免连接池成为新瓶颈。5. 从本地调试到生产部署三阶段验证法与MySQL写入性能压测技巧这套方案的价值不在“能跑”而在“能稳、能查、能扩”。我把它拆成三个递进阶段来验证——每个阶段都有明确的验收指标和失败回滚点避免陷入“改一点、试半天、不知哪错了”的泥潭。5.1 阶段一单机闭环验证15分钟内完成目标确认数据从Kafka到MySQL的端到端链路畅通且时间语义正确。操作清单启动ZooKeeper./zookeeper-3.4.11/bin/zkServer.sh start启动Kafka./kafka_2.10-0.9.0.0/bin/kafka-server-start.sh ./kafka_2.10-0.9.0.0/config/server.properties创建Topic./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic student-topic启动MySQL确保Student.sql已执行编译并运行Flink Jobmvn clean package flink run -c com.example.flink.kafka2mysql.KafkaToMysqlJob target/kafkasink2mysql-1.0.jar发送测试数据模拟学生事件echo {studentId:S001,score:85,event_time:2023-10-01 12:00:01} | \ ./kafka-console-producer.sh --broker-list localhost:9092 --topic student-topic验证指标Flink Web UI →Task Managers → Metrics → numRecordsInPerSecond 0MySQL执行SELECT * FROM student_agg WHERE idS001;确认total_score85,count1,last_update_time为当前时间误差2秒发送第二条数据{studentId:S001,score:92,...}再次查询确认total_score177,count2。注意若第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条用t0后250条用t10秒后模拟严重乱序 if [ $i -le 250 ]; then ts2023-10-01 12:00:00 else ts2023-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 idS001;预期结果count500total_score为500个随机数之和。若count500说明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 MEMORY8405006GSELECT COUNT(*) FROM information_schema.PROCESSLIST;压测执行步骤修改KafkaToMysqlJob.java中env.setParallelism(4)按上表调整MySQL配置并重启使用sysbench模拟写入压力避免干扰Flinksysbench --db-drivermysql --mysql-hostlocalhost --mysql-port3306 \ --mysql-userroot --mysql-passwordpassword --mysql-dbtest \ oltp_write_only --tables1 --table-size1000000 --threads32 run观察Flink Web UI的latency指标应500ms和MySQL的Threads_running应50。从那以后我每次上线新的Flink-Kafka-Mysql链路都强制走一遍这三阶段先用单条数据打穿链路再用乱序数据校验时间语义最后用sysbench压测DB水位。少走一次就多一个半夜被PagerDuty叫醒的理由。希望帮到你。本文还有配套的精品资源点击获取
返回列表