这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 是实时数据管道的事实标准。本章讲它与 Flink、Spark、ClickHouse、数据湖、CDC 的集成方式与关键取舍。
25.1 整体链路
业务库 --CDC(Debezium)--> Kafka --Flink/Spark-->
├-> 实时看板(ClickHouse/Doris)
├-> 数据湖(Iceberg/Hudi)
└-> Kafka(回流/宽表)
Kafka 在其中的职责:缓冲、解耦、回放、多订阅。一旦数据进了 Kafka,就可以被任意多个引擎按需消费,互不影响。
25.2 Kafka + Flink
Flink 是流处理生态中事件时间、状态管理与 exactly-once 能力最强的引擎之一,Kafka 是它最常见的 source/sink。
Maven 依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.19.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>3.1.0-1.19</version>
</dependency>
典型作业
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30_000);
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("broker1:9092")
.setTopics("user-behaviors")
.setGroupId("flink-behavior-job")
.setStartingOffsets(OffsetsInitializer.committedOffsets(
OffsetResetStrategy.EARLIEST))
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<Event> events = env.fromSource(
source, WatermarkStrategy
.forBoundedOutOfOrderness(Duration.ofSeconds(5)),
"kafka-source")
.map(this::parseEvent);
events.keyBy(e -> e.userId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new SessionAgg())
.sinkTo(buildClickHouseSink());
env.execute("behavior-agg");
Exactly Once 的两条线
1. Kafka -> Flink 内部:
checkpoint 把消费 offset 记入状态,故障恢复时从 checkpoint 位点重放
2. Flink -> Kafka 输出:
KafkaSink 配置 DELIVERY_GUARANTEE_EXACTLY_ONCE,
依赖 Kafka 事务 + 两阶段提交
KafkaSink 示例:
KafkaSink<AggResult> sink = KafkaSink.<AggResult>builder()
.setBootstrapServers("broker1:9092")
.setRecordSerializer(new SimpleStringSchema(), "result-topic")
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("flink-agg-")
.build();
注意:
- exactly-once sink 会拉长端到端延迟(受 checkpoint 与事务提交影响);
- 下游 Kafka 消费者要用
read_committed; - watermark 决定事件时间语义,乱序窗口要配 allowed lateness 与侧输出。
25.3 Kafka + Spark Structured Streaming
已有 Spark 技术栈时,Structured Streaming 也能消费 Kafka:
val df = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "broker1:9092")
.option("subscribe", "user-behaviors")
.option("startingOffsets", "latest")
.load()
val events = df.selectExpr(
"CAST(key AS STRING) key",
"CAST(value AS STRING) json"
).select(from_json($"json", schema).as("e")).select("e.*")
val agg = events
.withWatermark("event_time", "10 minutes")
.groupBy(window($"event_time", "5 minutes"), $"itemId")
.count()
agg.writeStream
.format("console")
.outputMode("update")
.start()
.awaitTermination()
对比:
| 维度 | Flink | Spark Streaming |
|---|---|---|
| 延迟 | 毫秒-秒级 | 微批(默认百毫秒-秒级) |
| 状态 | RocksDB 大状态成熟 | 基于 checkpoint |
| 事件时间/watermark | 强 | 支持 |
| 生态 | 纯流 | 批流一体、队列生态 |
25.4 Kafka -> ClickHouse / Doris
实时看板常用组合。三种摄入方式:
| 方式 | 特点 |
|---|---|
| 自研消费者批量 INSERT | 灵活、可控,需处理幂等与批次 |
| ClickHouse Kafka 表引擎 | 数据库内拉取,部署简单,运维边界模糊 |
| Doris/StarRocks Routine Load | 声明式任务,原生支持 group 与重试 |
ClickHouse 批量插入建议:
- 单批 1 万-10 万行或 10-100MB;
- 按
ORDER BY键排序写入,减少后台 merge; - 用 ReplacingMergeTree / 协处理器处理迟到重复;
- 消费者按分区并行,写库失败不提交位移。
25.5 Kafka -> 数据湖(Iceberg)
近实时数仓路径:
Kafka -> Flink Iceberg Sink ->
ODS 表 -> 批/流加工 -> DWD/DWS -> 查询引擎(Trino/Spark)
Flink 写 Iceberg 的关键点:
- exactly-once 依赖 checkpoint 提交快照(不同版本实现细节不同);
- 小文件问题:调大 checkpoint 间隔与文件滚动条件,定期 compaction;
- 到达时间 vs 事件时间:湖表分区按事件时间列,避免迟到数据写错分区;
- 回溯重建:Kafka 保留期内可重放,超过则需从湖表或归档补。
25.6 CDC:Debezium
CDC 把数据库变更变成流:
MySQL binlog -> Debezium(Connect) -> Kafka topic(每表一个)
价值:
- 缓存失效、搜索索引更新准实时;
- 数仓同步替代双写与定时拉取;
- 审计流水(完整变更历史)。
关键实践:
- topic 按
server.schema.table命名(默认规则),key 为主键; - snapshot 模式选择
initial/schema-only/never; - Debezium offset 与 schema history topic 必须高可靠(RF=3);
- 下游用 upsert 语义(Flink
upsert-kafkaconnector、CH ReplacingMergeTree); - 敏感列在 connector 配置里脱敏或加密。
25.7 反压与资源
大数据链路中常见“Kafka lag 告警,但 Broker 很闲”:
排查顺序:
1. 计算引擎 slot/并行度是否不足
2. 算子是否有数据倾斜(keyBy 热点)
3. 外部 sink 是否慢(DB 写入、湖提交)
4. checkpoint 是否过大过慢,拖住作业
处理:
- 并行度对齐 Kafka 分区数(source 并行 <= 分区数);
- 热点 key 加盐或预聚合;
- sink 攒批/异步化;
- 状态清理(TTL)避免 checkpoint 膨胀。
本章小结
- Kafka 是大数据链路的缓冲与回放层,多引擎共享同一份流;
- Flink 与 Kafka 组合可做到端到端 exactly-once,注意 checkpoint、事务与 read_committed;
- Spark 适合批流一体场景,纯低延迟大状态优先 Flink;
- ClickHouse/Doris 消费要攒批,湖表要管理小文件与分区;
- CDC 是数仓与缓存同步的正确姿势,offset/schema 主题必须高可靠。
思考题
- Flink checkpoint 间隔从 30s 调到 5 分钟,对吞吐、延迟、恢复各有什么影响?
- 为什么 Kafka -> Flink -> Kafka 的 exactly-once 在链路外还要
read_committed? - 设计一个“MySQL 订单表实时同步到 ClickHouse”方案,说明迟到与重复如何处理。