简介:一份基于 Apache Flink 的实时数据处理工程包,完整演示从 Kafka 消费数据、执行流式计算、再将结果写入 Redis 集群与 MySQL 的链路。面向大数据开发初学者及有实时数仓落地需求的工程师,适合用于学习 Flink Connector 配置、算子应用和外部存储集成。压缩包大小约 48.47MB,共 145 个文件,以 XML 配置、Java 源码、Class 编译结果和 Properties 配置为主,另有 Maven 构建脚本、依赖 Jar 及工程辅助文件,模块划分清楚。目前已有 597 人学习下载。通过该工程可掌握 Kafka 消费者组与分区消费配置、KeyBy 分组和 Window 窗口计算,以及数据写入 Redis 时键值或哈希结构的使用;同时可参考基于 JDBC 的 MySQL 批量写入实现,理解事务与一致性处理。由于工程还涉及 Redis 集群的槽位路由与故障恢复,对构建实时监控、在线广告分析等低延迟场景有直接参考价值。
1. flink读取kafka数据.zip:先拆开看它到底想让你做什么
如果你手里也有一份 flink读取kafka数据.zip,大概率是从博客、培训课件或同事网盘里扒下来的 demo 工程。它的目标很集中:把 Kafka topic 里的消息接进 Flink,实时做清洗、聚合、告警,再交给下游。适合两类人——刚接触实时计算、想在本机把整条链路跑起来的新手;已经能写简单 Flink 作业、但一上生产就遇到读不到数据、消息延迟高的工程师。别因为名字带 zip 就小看它,启动方式、依赖版本、checkpoint 策略这三样但凡错一个,本地跑得欢的代码,换台机器可能就静默不消费,比直接报错还难查。
2. 搭最小环境:Kafka 集群安装与 Flink 依赖选型
读 Kafka 之前必须先有一个能连的 Kafka,这一步很多人直接卡住。Kafka 集群安装 的官方文档写得很完整,但本地跑 demo 不需要三节点,下载官方 tar 包起单 broker 就够。常见做法是下载 Kafka 2.13 或 3.x 的二进制包,解压后一条命令起 ZooKeeper,另一条起 broker。别嫌老,ZooKeeper 模式依旧是排查问题最直接的模式。
2.1 先装 Kafka 集群:单机版足够跑通 demo 的安装步骤
Kafka 软件包内置了 ZooKeeper,本地验证连额外组件都不用装。先下载并解压,把KAFKA_HOME指到解压目录:
wget https://downloads.apache.org/kafka/3.6.1/kafka_2.13-3.6.1.tgz tar -xzf kafka_2.13-3.6.1.tgz cd kafka_2.13-3.6.1启动顺序固定:先 ZooKeeper,再 Kafka broker。两个窗口分开跑,因为一个终端占住进程后不方便看另一个日志。
bin/zookeeper-server-start.sh config/zookeeper.properties bin/kafka-server-start.sh config/server.properties这几条命令在 Linux 和 macOS 上直接能用。Windows 用户记得把bin换成bin\windows,后缀从.sh变成.bat,例如bin\windows\zookeeper-server-start.bat config\zookeeper.properties。有一个点经常被忽略——server.properties里listeners默认绑定0.0.0.0:9092,但如果不开advertised.listeners,broker 对外广播的地址是localhost:9092。你在本机连没事,另一台机器上的 Flink 作业连过来就会看到连接被重置或超时。
Kafka 起来以后建一个测试 topic。分区数先设成 3,后面讲并行度时会发现这个数字不是随便选的。
bin/kafka-topics.sh --create --topic demo-topic \ --bootstrap-server 127.0.0.1:9092 \ --partitions 3 --replication-factor 1单机版replication-factor只能填 1。partitions=3意味着这个 topic 最多能让 3 个 Flink source 子任务并行消费,后面 4.2 会展开。
2.2 Flink 连 Kafka 的依赖怎么选:版本对不上就是第一道坑
Kafka 本身跑通只完成了三分之一,Flink 工程里的依赖才是最常见的翻车现场。不同 Flink 版本对应的连接器坐标不一样,尤其是 1.15 之后连接器独立发布,老博客里的写法经常是过时的。下面是我常用的选择逻辑:
| Flink 版本系列 | 连接器写法 | 常用 API |
|---|---|---|
| 1.12 ~ 1.14 | flink-connector-kafka_2.12 | FlinkKafkaConsumer |
| 1.15 ~ 1.17 | flink-connector-kafka独立版本 | KafkaSource为主 |
| 1.18 及以上 | flink-connector-kafka | KafkaSource |
在 Maven 工程里至少要加入flink-streaming-java和flink-connector-kafka两个依赖。以 Flink 1.17 为例,pom 里写:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.2</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>3.0.2-1.17</version> </dependency>flink-streaming-java加provided是因为线上集群自带 Flink 依赖,不打进 jar 能大幅减小包体。但本地 IDEA 直接运行时,部分 IDEA 版本会把 provided 依赖排除出运行 classpath,导致启动就报NoClassDefFoundError。我一般本地调试时先把 provided 去掉,提交集群前再加回去,省得在环境问题上浪费时间。
连接器包的 scala 版本_2.12、_2.13也要和 Flink 发行版对应。混用会出现一堆NoSuchMethodError,而且报错位置千奇百怪,看起来跟 Kafka 无关,实际就是二进制不兼容。
2.3 自检连通性:topic、控制台消费、端口检查
写完代码之前,先用 Kafka 自带命令确认环境是通的。开两个终端,一个跑生产者,一个跑消费者:
bin/kafka-console-producer.sh --broker-list 127.0.0.1:9092 --topic demo-topic bin/kafka-console-consumer.sh --bootstrap-server 127.0.0.1:9092 --topic demo-topic --from-beginning在生产者窗口输入一行hello-kafka,消费者窗口能看到同样内容,说明 Kafka 自身没问题。这个自检很重要,它把问题边界划清楚了:之后 Flink 读不到数据,就不是 Kafka 的问题,而是 Flink 配置的问题。
再检查一下端口监听。Linux 上用ss -lntp | grep 9092,Windows 上用netstat -an | findstr 9092。如果端口没起来,回头看 broker 日志,多半是 ZooKeeper 没起或者config/server.properties里的log.dirs目录没有写权限。环境通了之后,建议装一个 Kafka 可视化工具,常见的有 kafka-ui 和 Offset Explorer。它们能直接看到 topic 分区、offset 和 consumer lag,后面 4.4 排查消息延迟时是刚需。这也回答了很多人问的“Kafka 有没有 UI 界面”——官方确实不附带,但第三方可视化工具很成熟,几十 MB 的安装包就能跑。
3. 用 Java 写出第一个 flink 读取 kafka 的最小程序
环境就绪后,写最小可跑的程序。这里我不会给你贴一个几十行的“完整工程”,而是按一个 Flink 作业该有的结构拆开讲:source 怎么建、反序列化怎么做、本地怎么跑。把这三块看懂,zip 里那个 demo 的代码你一眼就能读明白。
3.1 最小骨架:KafkaSource 与一个 print 算子
Flink 1.15 之后官方推荐用KafkaSource,它比老的FlinkKafkaConsumer干净,builder 模式把参数都显式暴露出来。一个能从 Kafka 读字符串并打印到控制台的最小作业长这样:
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.serialization.SimpleStringSchema; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class KafkaDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(30_000); KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("127.0.0.1:9092") .setTopics("demo-topic") .setGroupId("flink-demo-group") .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(); DataStream<String> lines = env.fromSource( source, WatermarkStrategy.noWatermarks(), "kafka-demo-source"); lines.print(); env.execute("flink-kafka-demo"); } }这段代码的逻辑分三层。第一层StreamExecutionEnvironment是作业的运行时容器,enableCheckpointing(30_000)表示每 30 秒做一次 checkpoint,是把 Kafka offset 提交和状态恢复绑定的前提。第二层KafkaSource负责和 Kafka broker 建立连接,setValueOnlyDeserializer表示只用 value 反序列化,key 忽略。第三层fromSource把 source 转成DataStream,print()把每条消息打到标准输出,env.execute()提交作业。
参数说明里最容易被忽略的是setStartingOffsets(OffsetsInitializer.earliest())。earliest表示新 group 从 topic 最老的 offset 开始读,latest表示只读启动之后新到的数据。如果你用默认的 latest,topic 里已经存在的消息一条都不会出现。调试阶段我建议显式指定earliest,上线时再根据业务改成latest或OffsetsInitializer.committedOffsets()。
老项目里常见的FlinkKafkaConsumer写法是env.addSource(new FlinkKafkaConsumer<>("demo-topic", new SimpleStringSchema(), props)),它确实还在跑,但从 Flink 1.15 起被标记为过时,新代码不建议再用。zip 里的 demo 如果是老写法,别急着换,功能等价,只是 API 风格不同。
3.2 反序列化与事件时间:字符串之外的选择
SimpleStringSchema只适合 Kafka 里存的是普通 UTF-8 文本。一旦 value 是 JSON,你可以先按字符串读进来再用JSON.parseObject去解析,但这会让类型信息在流里流动时丢失,后面做窗口聚合会很别扭。更干净的做法是实现一个反序列化器,让 Flink 直接产出 POJO。下面这个例子把字节转成Event对象:
import org.apache.flink.api.common.serialization.AbstractDeserializationSchema; import com.fasterxml.jackson.databind.ObjectMapper; import java.io.IOException; public class EventDeserializer extends AbstractDeserializationSchema<Event> { private static final ObjectMapper MAPPER = new ObjectMapper(); @Override public Event deserialize(byte[] message) throws IOException { return MAPPER.readValue(message, Event.class); } }用的时候把 source 部分的setValueOnlyDeserializer(new SimpleStringSchema())换成setValueOnlyDeserializer(new EventDeserializer())即可。注意这个类一定要写成独立类,不能写成匿名内部类。Flink 的TypeInformation推断在匿名内部类上会拿不到泛型信息,运行时报一堆类型异常。这个坑我踩过,印象极深。
如果后续要按事件时间做窗口,比如统计每分钟的告警数量,必须在 source 之后指定 timestamp 和 watermark:
DataStream<Event> eventsWithWatermark = lines .assignTimestampsAndWatermarks( WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, ts) -> event.getTs().getTime()));forBoundedOutOfOrderness(Duration.ofSeconds(5))表示允许事件时间最多乱序 5 秒。这个参数不是随便填的,它取决于上游 Kafka 消息的产生机制。如果业务方有多个服务写同一个 topic,时钟偏差可能超过 10 秒,watermark 填太小会导致大量窗口数据晚到被丢弃,填太大又会延迟窗口触发。先按 5 秒起步,观察延迟率再调。
3.3 本地提交:IDEA 里直接跑与并行度的关系
在 IDEA 里直接右键跑 main 方法,Flink 默认并行度等于你电脑的 CPU 核数。假设你的机器是 8 核,source 会起 8 个子任务,但 demo-topic 只有 3 个分区。多出来的 5 个子任务一直处于空闲状态,看起来像卡住,实际上是在空转等分区分配。
本地验证阶段,我建议在代码里加一行env.setParallelism(3),和 topic 分区数保持严格一致。这样做的原因是让并行度暂态和分区数对齐,后续调参时不容易被空转子任务干扰。生产上并行度可以超过分区数,但超过部分没有收益,白白占 slot。
还有打包问题。如果你想把作业提交到 Standalone 集群,pom 里要用maven-shade-plugin把连接器依赖打进去,否则提交后报NoClassDefFoundError: org/apache/flink/connector/kafka/source/KafkaSource。shade 插件配置很固定,核心是createDependencyReducedPom设为 false,让所有依赖进入 fat jar。本地 IDEA 跑不需要打包,直接 run 就行。
4. 生产要调的 Kafka 消费参数:并行度、offset、语义
从“能跑”到“敢上生产”,中间隔着几个必调参数。很多人觉得能从 Kafka 读到数据就够了,结果第一次重启作业就丢数据或重复消费。这一章专门讲参数怎么设,以及每条参数背后的理由。
4.1 group.id 与 offset 的起点:重启后不丢数据的第一道防线
Kafka 的 offset 天然存在内置的__consumer_offsetstopic 里,但 Flink connector 并不是主要提交者——Flink 只在 checkpoint 完成时把 offset 提交给 Kafka。这个差异特别重要:如果你在 Flink 作业里开enable.auto.commit=true,Flink 会直接禁用自动提交。理解这一点,后续调参才不会懵。
客户端连接参数通过Properties透传给 Kafka 生产者,但我更推荐直接在KafkaSource.builder()上设置。几个核心参数可以整理成一张表:
| 参数 | 默认值 | 作用 | 我的建议 |
|---|---|---|---|
bootstrap.servers | 无 | broker 地址列表 | 写多台,逗号分隔 |
group.id | 无 | 消费组,决定 offset 归属 | 按业务线命名,不要随便改 |
auto.offset.reset | latest | 新 group 无 offset 时的起点 | 调试用 earliest,生产用 committedOffsets |
enable.auto.commit | true | 是否自动提交 offset | 交给 Flink checkpoint,不必显式关 |
改group.id要非常谨慎。Flink 检查点里记着当前的消费位点,如果你换一个 group.id 再恢复,等于让作业从“失忆”状态重新开始。我见过有人每次发布改 group.id 做灰度,结果所有任务从头消费,Kafka 压测直接打满。
4.2 并行度与分区数:不是越大越好
并行度和 Kafka 分区数的关系,是理解 source 性能的关键。Flink 的 Kafka source 每个子任务消费一个或多个分区,但一个分区同一时刻只能被一个子任务消费。分区数和 source 并行度匹配时有三种情况:
分区数等于并行度时,每个子任务恰好消费一个分区,负载最均匀。分区数小于并行度时,部分子任务空转,不读数据但占着 slot。分区数大于并行度时,一个子任务串行消费多个分区,单个分区如果消息量大,这个子任务就成为瓶颈。
所以当 kafka 消息延迟高 的时候,第一反应不应该是把并行度翻倍,而要先看 topic 分区数够不够。我曾经把一个 3 分区的 topic 配了 12 并行度,数据量上来后延迟照旧,因为 3 个分区只能撑起 3 个有效消费线程。后来把分区扩到 12,问题立刻缓解。分区数扩容在 Kafka 侧一条命令就能做,但要注意扩容后 key 在分区之间的映射会变,如果下游按分区有序消费,要先确认业务能不能接受。
4.3 checkpoint 和消费语义:至少一次还是精确一次
Flink 的消费语义完全建立在 checkpoint 之上。没有 checkpoint,Flink 就退化成“Kafka 客户端自动提交 + 内存处理”,重启后要么重复要么丢失,没有任何保障。生产环境必开 checkpoint,下面是我常用的配置:
env.enableCheckpointing(10_000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5_000); env.getCheckpointConfig().setCheckpointTimeout(60_000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);enableCheckpointing(10_000)是每 10 秒做一次快照,间隔太短会频繁访问状态后端,太长则故障恢复时丢失的数据多。setMinPauseBetweenCheckpoints(5000)保证两个 checkpoint 之间至少间隔 5 秒,防止背压时 checkpoint 连环失败。EXACTLY_ONCE模式配合 Kafka source 和事务型 sink,能做到端到端精确一次,但代价是 sink 要等 checkpoint 完成才提交数据,延迟会略高。大多数报表场景用AT_LEAST_ONCE就够,精确一次别盲目开。
4.4 消息延迟高的排查路径:背压、lag 与 1M 消息限制
kafka 消息延迟高不是一个单一原因,而是一条链路的多米诺骨牌。我的排查顺序固定不变:先看消费者 lag,再看 Flink UI 的背压,最后看消息大小和下游写入速度。
第一步看 lag。用kafka-consumer-groups.sh或可视化工具查 group 的CurrentOffset和LogEndOffset,lag 持续增长说明消费速度跟不上生产速度。第二步进 Flink UI,看 source 上游的算子是不是BackPressure HIGH。如果背压出现在 source 上,说明不是 Kafka 读得慢,而是下游算子处理不动。第三步看 sink。很多人只盯着 source 调并行度,结果瓶颈在 JDBC 写入,source 并行度再高也没用,数据全堵在 sink 前面的队列里。
还有一个和“kafka 接收 1m”相关的坑:Kafka broker 默认message.max.bytes=1MB,consumer 端的fetch.max.bytes也有默认上限。如果业务确实要传 1MB 级的大消息,光调 broker 端不够,Flink 连接器侧的fetch.max.bytes也要同步调大,否则源头就丢不进。调大后还要注意内存占用,每条 1MB 的消息在算子间传递拷贝几次,GC 会变成新的瓶颈。
5. flink 读取 kafka 的避坑清单:五条翻车记录
下面这些是我在不同项目里实际遇到过的坑,每一条都按“现象、原因、解决”写成了一条记录。前四条在本地就能复现,最后一条只发生在跨机器部署时。你如果正在调试,按序号对照排查,比从头看日志快得多。
5.1 子任务卡在 INITIALIZING:并行度大于分区数导致空转
现象:作业提交后,Flink UI 里 source 子任务长时间停在INITIALIZING或STARTING,日志里只有几条 Kafka 连接记录,没有消费数据。
原因:最常见的是并行度大于 topic 分区数。多余的子任务拿不到分区,一直在等待,状态就停在启动阶段。还有可能是bootstrap.servers写成了域名,但本机没配 hosts,导致连接请求超时后反复重试。
解决:先kafka-topics.sh --describe --topic demo-topic看分区数,再检查代码里env.setParallelism(1)或 source 的setParallelism(3)是否和分区数匹配。我习惯把setBootstrapServers写成 IP:9092 的形式,避免 DNS 解析这一层不确定性。
5.2 作业起来了但 print 一直没输出:新 group 从 latest 开始读
现象:Flink 作业正常 submit,print()也执行了,但控制台什么都没有。往 Kafka 里再发新消息,突然又有输出了。
原因:auto.offset.reset默认latest,新 group 首次消费时没有已提交 offset,会从连接建立之后的“最新位置”开始读。如果先启动 Flink,再往 Kafka 里发消息,那是对的;但很多调试场景是 Kafka 里已经有积压消息,再启动 Flink,旧消息就全部被跳过。
解决:调试阶段用OffsetsInitializer.earliest()指定从最早读。生产上如果首次启动想消费积压数据,同样显式配earliest;如果只处理增量数据,latest也可以,但必须知道这是有意为之。
5.3 反序列化报错:SimpleStringSchema 不是万能的
现象:deserialize抛JsonParseException或ClassCastException,日志里显示一堆字节转字符串后乱码。
原因:SimpleStringSchema把 Kafka value 按 UTF-8 字节直接转 String,如果生产者用的是 Avro、Protobuf 或自定义二进制格式,字符串里就会出现乱码,继续 JSON 解析必然失败。还有可能是同一个 topic 里混了多种消息格式,这在跨团队协作的 topic 上很常见。
解决:先确认生产者序列化格式,再写对应的AbstractDeserializationSchema。如果 topic 混合多格式,就别在 source 层面强行转换,先把消息按byte[]读进来,再根据消息头字段分发到不同解析逻辑。这个解耦方案虽然要写更多代码,但能避免后续接新格式时反复改动 source。
5.4 重启作业后数据大量重复:checkpoint 没开,offset 全凭自动提交
现象:作业崩溃后重启,从很久以前的 offset 开始消费,数据重复量巨大;或者相反,重启后直接跳过了一部分数据。
原因:没有enableCheckpointing,Flink 不持有 offset 的控制权,届时用 Kafka 客户端的enable.auto.commit按默认 5 秒间隔自动提交。作业在两次自动提交之间崩溃,重启时 offset 回退到上次提交点,这个窗口内的消息全部被重复消费。
解决:开 checkpoint,并把checkpointingMode设为EXACTLY_ONCE或AT_LEAST_ONCE。Flink 会在每次 checkpoint 完成时把消费位点提交到 Kafka,恢复时从最近一次成功的 checkpoint 接着读。这里要注意:checkpoint 频次定义了你最多可能重复多少数据。10 秒一次 checkpoint,崩溃时最多重复最近 10 秒的数据。要更小的重复窗口,就把 checkpoint 间隔缩到 3 秒,但状态后端压力会明显增大。
5.5 Windows 本地连 Kafka 超时:advertised.listeners 背锅
现象:IDEA 里跑 Flink 作业,一两分钟后报TimeoutException,但同一个 Kafka 在其它机器上消费正常。
原因:broker 的advertised.listeners没有显式配置,默认广播localhost:9092。Windows 本机连的时候没事,但你的 Flink 作业和 Kafka 之间只要有网络转发,客户端拿到的 broker 地址就是 localhost,根本路由不过去。
解决:在config/server.properties里加显式配置:
listeners=PLAINTEXT://0.0.0.0:9092 advertised.listeners=PLAINTEXT://你的本机IP:9092改完一定要重启 Kafka broker 才生效。如果改了还超时,Windows 防火墙把 9092 端口拦了再加一条入站规则。判断到底是不是这个问题的技巧:在 Flink 作业所在机器上用telnet 你的IP 9092,如果能通但 Flink 仍然超时,再往连接器配置和 zookeeper 连接方向排查。
6. 进阶:从读 Kafka 到写 ClickHouse,做一条完整的本地链路
读 Kafka 只是第一步,落地到存储才算把链路打通。很多人在这个环节卡在 flink 的 jdbc 连接器异常上,这里用 ClickHouse 做例子,因为它的 JDBC 接入方式和 MySQL 风格不太一样,踩坑率高,一次性讲透。
6.1 再挂一个 JdbcSink:把结果落到 ClickHouse
给 Kafka 里每个Event对象配一个JdbcSink,按事件时间写入 ClickHouse 表。这里直接复用上文的EventDeserializer和eventsWithWatermark流:
import org.apache.flink.connector.jdbc.JdbcExecutionOptions; import org.apache.flink.connector.jdbc.JdbcSink; import org.apache.flink.connector.jdbc.JdbcConnectionOptions; eventsWithWatermark.addSink( JdbcSink.sink( "INSERT INTO demo.events(ts, value) VALUES (?, ?)", (ps, event) -> { ps.setTimestamp(1, new Timestamp(event.getTs().getTime())); ps.setString(2, event.getValue()); }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchIntervalMs(5000) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl("jdbc:clickhouse://127.0.0.1:8123/demo") .withDriverName("com.clickhouse.jdbc.ClickHouseDriver") .build()));逻辑上,JdbcSink.sink的参数分三段:SQL 语句、每行参数绑定逻辑、连接配置。withBatchSize(1000)是一批攒 1000 条再写,withBatchIntervalMs(5000)是兜底机制,防止低流量时数据一直攒着不写。这两个参数一个控制吞吐,一个控制延迟上限。
ClickHouse 的连接会有两个特有问题,都是 flink 的 jdbc 连接器异常 里常见报错。第一个是No suitable driver,原因是驱动类名写成了ru.yandex.clickhouse.ClickHouseDriver老驱动,换成com.clickhouse.jdbc.ClickHouseDriver并引入新版clickhouse-jdbc依赖即可。第二个是建表时没指定引擎或 ORDER BY,ClickHouse 会拒绝写入。先建表再启 Flink 作业:
clickhouse-client --query " CREATE TABLE IF NOT EXISTS demo.events ( ts DateTime, value String ) ENGINE = MergeTree() ORDER BY ts"6.2 本地全链路验证:从 console producer 到 ClickHouse 行数
完整链路验证我用四步。第一步,启动 Kafka、ClickHouse、Flink 作业。第二步,用 Kafka console producer 发送 10 条 JSON 格式的 Event 消息,比如{"ts":"2024-06-01 12:00:00","value":"test-1"}这类格式,注意和Event里的字段对上。第三步,到 Flink 控制台看有没有 10 行 print 输出,再到 ClickHouse 执行SELECT count(*) FROM demo.events,确认入库条数。第四步,也是最硬核的一步——验证 checkpoint 是否真的在起作用:让作业运行一会,直接 kill -9 进程,再重新启动,观察 ClickHouse 里的数据是否重复。
第四步能暴露很多日志看不到的问题。如果没有重复,说明 checkpoint 已经把位点和 sink 事务绑在了一起;如果重复,说明 checkpoint 频率不够或者 sink 配置了幂等之外的行为。这套验证做完,作业才有资格往测试环境推。
我自己的习惯是,任何 Flink 流任务提交之前,一定先把这条本地链路完整跑一遍:Kafka 生产、Flink 消费、ClickHouse 落库、kill 重启验重复。延迟、重复、断流,这些生产上最难查的问题,其实都能在本地用几分钟模拟出来。别等生产环境帮你测,到那时你已经没有干净的日志可以看了。希望帮到你。
本文还有配套的精品资源,点击获取