KafkaNotes

第 25 章:Kafka 与大数据生态

zjc 于 2026-01-25 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 是实时数据管道的事实标准。本章讲它与 Flink、Spark、ClickHouse、数据湖、CDC 的集成方式与关键取舍。

25.1 整体链路

业务库 --CDC(Debezium)--> Kafka --Flink/Spark-->
  ├-> 实时看板(ClickHouse/Doris)
  ├-> 数据湖(Iceberg/Hudi)
  └-> Kafka(回流/宽表)

Kafka 在其中的职责:缓冲、解耦、回放、多订阅。一旦数据进了 Kafka,就可以被任意多个引擎按需消费,互不影响。

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();

注意:

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 批量插入建议:

25.5 Kafka -> 数据湖(Iceberg)

近实时数仓路径:

Kafka -> Flink Iceberg Sink ->
  ODS 表 -> 批/流加工 -> DWD/DWS -> 查询引擎(Trino/Spark)

Flink 写 Iceberg 的关键点:

25.6 CDC:Debezium

CDC 把数据库变更变成流:

MySQL binlog -> Debezium(Connect) -> Kafka topic(每表一个)

价值:

关键实践:

25.7 反压与资源

大数据链路中常见“Kafka lag 告警,但 Broker 很闲”:

排查顺序:
1. 计算引擎 slot/并行度是否不足
2. 算子是否有数据倾斜(keyBy 热点)
3. 外部 sink 是否慢(DB 写入、湖提交)
4. checkpoint 是否过大过慢,拖住作业

处理:

本章小结

思考题

  1. Flink checkpoint 间隔从 30s 调到 5 分钟,对吞吐、延迟、恢复各有什么影响?
  2. 为什么 Kafka -> Flink -> Kafka 的 exactly-once 在链路外还要 read_committed
  3. 设计一个“MySQL 订单表实时同步到 ClickHouse”方案,说明迟到与重复如何处理。