Kafka 学习笔记

从内部原理到生产运维,把 Kafka 彻底讲透。不写安装教程,只讲架构设计、存储引擎、副本机制、性能调优和踩坑经验。

📌 学习路线建议

这篇笔记内容较多,建议分五步走:

阶段 读哪些章节 目标
第一步:理解全貌 架构篇 第 1~3 节 搞懂 Kafka 整体架构、核心概念、Controller 与 KRaft
第二步:深入存储 存储引擎篇 第 4~6 节 理解 Log Segment、索引、零拷贝、日志清理
第三步:生产者/消费者 第 7~10 节 分区策略、acks、幂等事务、Rebalance、Offset 管理
第四步:高可用 副本篇 第 11~13 节 ISR、Leader 选举、HW 与 Leader Epoch
第五步:生产加固 调优与运维篇 第 14~17 节 性能调优、监控、分区重分配、常见故障

建议结合消息队列学习笔记一起看——那篇讲 MQ 通用原理(为什么需要、消息丢失/重复/顺序),这篇讲 Kafka 专属的深度机制。


0 · 消息全生命周期

在深入每个组件之前,先跟着一条消息走一遍完整旅程。这一节只求建立直觉,不展开细节——每个步骤后面都有专门章节深入讲解。

图 5 \xb7 消息全生命周期

生产端:消息怎么发出去的

假设你是一个 Producer(生产者),要发一条订单消息到 Kafka。你会经历以下步骤:

第 1 步:序列化 —— 把你的订单对象转成字节数组。因为网络传输只认字节,不认你的 C++ struct 或 Java 对象。最简单的就是转成 JSON 字符串再编码成字节,复杂场景可以用 Avro 或 Protobuf(一种更紧凑的二进制格式,后面会讲)。

第 2 步:选分区 —— 一个 Topic(话题)可以分成多个 Partition(分区),就像一本书分成多册。你的消息要放进哪一册?有三种方式:

  • 你自己指定"放第 2 册"(直接指定分区号)
  • 消息带了一个 key(比如用户 ID),Kafka 对 key 做哈希运算,算出放哪册——好处是同一个用户的消息永远进同一分区,保证顺序
  • 消息没有 key,Kafka 帮你轮着来——第一条进分区 0,第二条进分区 1,依此类推,让消息均匀分布

第 3 步:攒批发送 —— Kafka 不是来一条发一条,而是先攒着。就像快递员不会拿到一个包裹就跑一趟,而是等包裹攒到一定数量,或者等了一定时间,才一起送出去。这样一次网络请求能带很多消息,效率高得多。

第 4 步:写入 Broker —— 消息到达 Broker(Kafka 服务器)后,Broker 里的"Leader"节点负责把消息追加写入一个日志文件。注意是追加,不是修改——永远只在文件末尾往后写,所以速度极快(顺序写磁盘的速度接近写内存)。

第 5 步:Follower 同步 —— 如果这个分区有副本(为了高可用),其他 Broker 上的"Follower"会主动从 Leader 那里拉取这条新消息,保持数据一致。注意是 Follower 主动拉,不是 Leader 推——这样 Follower 可以按自己的节奏同步,不会被压垮。

消费端:消息怎么被消费的

现在消息已经存在 Broker 里了,Consumer(消费者)要来取:

第 6 步:拉取消息 —— Consumer 主动向 Broker 发请求:"给我从第 N 条开始的消息"。这是 Pull 模型——消费者自己去拿,而不是 Broker 硬塞。好处是消费者可以按自己的速度拿,不会被压垮。

第 7 步:反序列化 —— Broker 返回的是字节数组,消费者要把它还原成订单对象。跟第 1 步的序列化是反操作——用什么格式序列化的,就得用什么格式反序列化。

第 8 步:业务处理 —— 消费者拿到订单对象后,执行业务逻辑:扣库存、发短信、更新积分等。这一步要注意幂等(同一条消息被消费多次也不能出问题),因为网络故障可能导致重复消费。

第 9 步:提交 Offset —— 处理完后,消费者告诉 Broker"我已经处理到第 N 条了"。这个位置标记就叫 Offset(偏移量)。下次启动时,消费者从上次记录的 Offset 继续消费,不会重复。如果消费者挂了没来得及提交,重启后会从旧 Offset 重新消费——这就是"At Least Once"语义,所以才需要幂等。

📌 一句话总结

生产端是 推 模型(Producer 主动推给 Broker),消费端是 拉 模型(Consumer 主动从 Broker 拉)。这条旅程中的每一步——序列化、分区、攒批、存储、同步、拉取、提交 Offset——后面都有专门章节展开讲。


第一部分 · 架构篇

这部分讲 Kafka 的整体架构和核心概念。读完你能搞懂 Kafka 的全貌:组件怎么协作、数据怎么流转、为什么这么设计。


1 · Kafka 整体架构

图 1 \xb7 Kafka 整体架构

1.1 核心组件

组件 职责 类比
Broker 存储和转发消息的服务节点 邮局的分拣中心
Producer 发送消息的客户端 寄信人
Consumer 消费消息的客户端 收信人
Controller 管理集群状态(分区/副本/Leader 选举) 邮局局长
ZooKeeper / KRaft 元数据存储 + Controller 选举 邮局档案室

1.2 一条消息的生命周期

图 5 \xb7 消息全生命周期

关键点:

  • Producer 推消息到 Broker(不是 Broker 拉)
  • Consumer 拉消息从 Broker(不是 Broker 推)
  • Follower 也是拉的——从 Leader 拉取副本数据
  • 消息写入只有 Leader 负责,Follower 只同步不写

1.3 为什么 Kafka 这么快

三个核心设计决策:

设计 原理 效果
顺序写磁盘 追加到文件末尾,不做随机写 磁盘顺序写速度 ≈ 内存随机写(600MB/s+)
页缓存(Page Cache) 消息先写 OS 页缓存,OS 异步刷盘 写入不走 JVM 堆,无 GC 问题
零拷贝(sendfile) 消费时数据从页缓存直接到网卡,跳过用户态 减少 2 次内核态↔用户态拷贝

一句话:Kafka 的快不是用了什么黑科技,而是把操作系统的能力用到了极致——顺序写、页缓存、零拷贝,三个 OS 级优化叠加,让磁盘 I/O 接近内存速度。


2 · 核心概念详解

2.1 Topic 与 Partition

  • Topic——逻辑分类,类似数据库的表。比如 orders、app-logs
  • Partition——Topic 的物理分片,是并行度的最小单位。每个 Partition 是一个有序的、不可变的追加日志
Topic: orders (3 partitions) ┌─────────────────────────────────────────────────────┐ │ Partition 0 │ Partition 1 │ Partition 2 │ │ [msg0, msg3, │ [msg1, msg4, │ [msg2, msg5, ...] │ │ msg6, ...] │ msg7, ...] │ │ └─────────────────────────────────────────────────────┘

为什么分区?

单机吞吐有上限。把一个 Topic 拆成 N 个 Partition,分布到 N 个 Broker 上,就能 N 倍并行写入和消费。分区数 = 并行度的上限。

分区数怎么定? 经验公式:分区数 = min(目标吞吐 / 单分区吞吐, Broker 数 × 消费者数)。一般起步 6~12 个分区,不够再加。注意:分区数只能加不能减。

2.2 Offset

每条消息在 Partition 中有一个单调递增的 64 位整数编号——Offset。

Partition 0: [0] [1] [2] [3] [4] [5] [6] [7] ... ↑ Consumer 当前位置
  • Offset 由 Broker 分配,不由 Producer 指定
  • Consumer 通过 Offset 控制消费位置——可以前进、可以回退、可以跳转
  • Offset 是 Kafka 简洁性的核心:不需要复杂的 ACK 机制,消费者只需要记住"我读到哪了"

2.3 消费者组(Consumer Group)

消费者组是 Kafka 实现发布订阅 + 点对点的统一机制:

  • 同一 Group 内:一个 Partition 只能被一个消费者消费 → 点对点
  • 不同 Group 间:各自独立消费全量消息 → 发布订阅
Topic: orders (3 partitions) Group A (订单处理): Group B (数据分析): Consumer 1 ← P0, P1 Consumer 4 ← P0, P1, P2 Consumer 2 ← P2 (只有1个消费者, 消费全部分区) Group C (通知服务): Consumer 5 ← P0 Consumer 6 ← P1 Consumer 7 ← P2 (3个消费者, 各消费1个分区)

📌 关键规则:一个 Group 内,分区数 = 最大消费者数。消费者数 > 分区数时,多余的消费者闲置。


3 · Controller 与 KRaft 模式

3.1 Controller 是什么

Controller 是 Broker 集群中的一个特殊角色——由其中一个 Broker 兼任。负责:

  • 分区 Leader 选举——Broker 宕机时,为该 Broker 上的 Leader 分区选出新 Leader
  • 分区创建/删除——管理 Topic 的生命周期
  • 副本重分配——扩容时迁移分区
  • Broker 上下线感知——监听 Broker 的心跳

3.2 ZooKeeper 模式(传统)

Kafka 2.x 及之前,依赖 ZooKeeper 存储:

元数据 存储位置
Broker 注册信息 /brokers/ids
Topic 配置 /brokers/topics
Partition 状态 /brokers/topics/{topic}/partitions
Controller 选举 /controller 临时节点
Consumer Offset(旧版) /consumers/{group}/offsets

ZooKeeper 的痛点:Kafka 依赖一个外部组件,运维复杂度翻倍。ZooKeeper 集群本身也要 3~5 台机器,元数据多了之后 Controller failover 慢(几十秒到分钟级)。

3.3 KRaft 模式(Kafka 3.x+,去 ZooKeeper 化)

KRaft(Kafka Raft)是 Kafka 内置的共识协议,替代 ZooKeeper:

图 6 \xb7 KRaft 模式架构

对比 ZooKeeper 模式 KRaft 模式
外部依赖 需要 ZooKeeper 集群 不需要,Kafka 自包含
元数据存储 ZooKeeper + Broker 内存 Broker 内部日志(Raft Log)
Controller failover 10~60 秒 1~3 秒
支持的分区数 ~20 万 百万级
运维复杂度 高(两套系统) 低(一套系统)

KRaft 的核心改进:元数据本身就是一个 Kafka Topic(__cluster_metadata),用 Raft 协议保证一致性。Controller 节点可以和 Broker 节点共用或分离(process.roles=controller,broker 或 process.roles=controller)。

# KRaft 模式配置示例 (docker-compose)
services:
  kafka:
    image: bitnami/kafka:3.7
    environment:
      KAFKA_CFG_NODE_ID: 0
      KAFKA_CFG_PROCESS_ROLES: controller,broker  # ★ 同时担任 Controller 和 Broker
      KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093  # ★ Raft 投票成员
      KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT

📌 趋势:Kafka 4.0 将完全移除 ZooKeeper 支持,KRaft 是未来。新项目直接用 KRaft 模式。


第二部分 · 存储引擎篇

这部分讲 Kafka 的消息是怎么存的、怎么找的、怎么读的。读完你能搞懂为什么 Kafka 能做到百万级 TPS。


4 · Log Segment 机制

4.1 Partition 的物理结构

每个 Partition 对应一个目录,里面存放多个 Log Segment(日志段):

/kafka-logs/orders-0/ ├── 00000000000000000000.log ← 第1个 Segment (起始 offset = 0) ├── 00000000000000000000.index ← 偏移量索引 ├── 00000000000000000000.timeindex ← 时间戳索引 ├── 00000000000000123456.log ← 第2个 Segment (起始 offset = 123456) ├── 00000000000000123456.index ├── 00000000000000123456.timeindex └── leader-epoch-checkpoint ← Leader Epoch 信息
  • 文件名 = 该 Segment 的起始 Offset,用 20 位数字补零
  • 只有一个 Segment 是 Active Segment(正在写入的),其余是已封存的
  • Active Segment 达到 log.segment.bytes(默认 1GB)或 log.roll.hours(默认 168 小时 = 7 天)后滚动为新 Segment

4.2 为什么用多 Segment 而不是单文件

原因 说明
日志清理 可以按 Segment 整体删除,不用逐条删
索引效率 每个 Segment 有独立索引,二分查找更快
并发读取 消费者读旧 Segment 时,不影响写入新 Segment

类比:如果整个 Partition 是一本书,Segment 就是章节。删除旧消息就是撕掉旧章节,不用一页页撕。写入新消息就是在最新章节末尾追加。


5 · 索引文件与零拷贝

图 3 \xb7 存储引擎结构

5.1 偏移量索引(.index)

索引文件是稀疏索引——不是每条消息都有索引项,而是每隔一定字节(index.interval.bytes,默认 4KB)记一条:

.index 文件内容: offset(相对) position(物理位置) 0 0 128 4096 ← 消费者要找 offset=130 256 8192 ← 二分查找: 130 在 128~256 之间 384 12288 → 跳到 position=4096, 顺序扫描到 130

查找流程:

  1. 二分查找 .index 文件,定位到 ≤ 目标 offset 的最近索引项
  2. 跳到 .log 文件对应的物理位置
  3. 顺序扫描 .log,直到找到目标 offset 的消息

为什么用稀疏索引而不是密集索引? 索引文件太小放不进所有消息的映射;太大会占用内存。稀疏索引让索引文件体积可控(通常几 MB),又能快速定位到附近,再顺序扫描几条就能找到。空间和时间的经典折中。

5.2 时间戳索引(.timeindex)

按时间戳查找消息的索引,用于"查找某个时间点之后的消息":

# 按时间戳查找 offset
kafka-run-class.sh kafka.tools.GetOffsetShell \
  --broker-list localhost:9092 --topic orders \
  --time 1704067200000  # 2024-01-01 00:00:00 的毫秒时间戳

5.3 零拷贝(sendfile)

传统数据读取(4 次拷贝 + 2 次系统调用):

磁盘 → 内核页缓存 → 用户空间缓冲区 → Socket 缓冲区 → 网卡 (内核态) (用户态) (内核态)

Kafka 的零拷贝(2 次拷贝 + 1 次系统调用):

磁盘 → 内核页缓存 → 网卡 (内核态) (DMA 直接到网卡)
// Kafka 底层用 Java 的 FileChannel.transferTo() → Linux sendfile()
// 消费者读消息时, Broker 直接从页缓存发到网卡, 不经过 JVM 堆

// C++ 客户端 (librdkafka) 消费端不需要关心零拷贝——这是 Broker 端的优化
// 你只需要正常 poll() 消费即可
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "my-group"},
    {"fetch.min.bytes", 1024},   // ★ 批量拉取, 减少网络往返
    {"fetch.max.wait.ms", 500},  // ★ 最多等 500ms 凑批量
});

零拷贝的前提:消息在页缓存中。如果消息不在页缓存(冷数据),还是要从磁盘读。所以 Kafka 的快依赖于热数据在页缓存——消费者跟得上生产者速度时,消息刚写入页缓存就被消费走,根本不落盘就被读走了。


6 · 日志清理策略

Kafka 有两种日志清理策略,可以按 Topic 配置:

6.1 按时间删除(默认)

# 默认保留 7 天
log.retention.hours=168

# 配置 Topic 只保留 1 小时
kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --topic orders \
  --add-config retention.ms=3600000

6.2 按大小删除

# 默认每个 Partition 保留 1GB
log.retention.bytes=1073741824

# 配置 Topic 每个 Partition 最多 500MB
kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --topic orders \
  --add-config retention.bytes=524288000

6.3 日志压缩(Log Compaction)

不同于直接删除,日志压缩是对同一个 Key 只保留最新值:

清理前 (Key → Value): user1 → {name: "Alice", age: 25} user2 → {name: "Bob", age: 30} user1 → {name: "Alice", age: 26} ← user1 更新了 user3 → {name: "Carol", age: 20} user1 → {name: "Alice", age: 27} ← user1 又更新了 清理后 (只保留每个 Key 的最新值): user2 → {name: "Bob", age: 30} user3 → {name: "Carol", age: 20} user1 → {name: "Alice", age: 27} ← 只留最新的
# 启用日志压缩
kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --topic user-states \
  --add-config cleanup.policy=compact   # ★ compact 而非 delete

适用场景:日志压缩适合"状态快照"类 Topic——比如用户的最新状态、配置的最新值。不适用于事件流(每个事件都是独立的,不能丢)。

delete + compact 可以共存:cleanup.policy=delete,compact —— 既按时间删除旧数据,又对保留的数据做压缩。


第三部分 · 生产者与消费者

这部分讲 Producer 和 Consumer 的内部机制和最佳实践。读完你能搞懂分区策略、acks 语义、幂等事务、Rebalance 原理。


7 · 生产者深入

7.1 分区策略

回忆第 0 节,一条消息发到 Topic 时要选择放进哪个 Partition。这一节展开讲三种选择方式。

为什么需要分区策略? 一个 Topic 有多个 Partition,消息进哪个 Partition 直接影响:

  • 负载均衡——消息均匀分布,每个 Broker 都有事干
  • 顺序保证——同一个 key 的消息进同一分区,分区内有序
  • 并行度——分区数决定了消费者最多能有多少个并行

策略一:直接指定分区

最简单的方式——你自己决定消息进哪个分区:

producer.produce(MessageBuilder("orders")
    .partition(2)         // ★ 强制发到 partition 2
    .payload(order_data)
);

适用场景:你需要精确控制,比如"VIP 用户的订单全部进分区 0,用专门的消费者处理"。但一般不推荐——硬编码分区号会导致负载不均。

策略二:有 Key 分区(最常用)

消息可以带一个 key(比如用户 ID、订单 ID)。Kafka 会对 key 做一次哈希运算(把任意长度的输入变成一个固定大小的数字),然后用这个数字对分区数取模,算出分区号:

分区号 = hash(key) % 分区数 例:hash("user-1001") = 7, 分区数 = 3 7 % 3 = 1 → 进入 Partition 1

关键特性:只要分区数不变,同一个 key 永远进同一个分区。这意味着同一个用户的所有消息是按顺序排列的——这对很多业务很重要(比如"先创建订单再取消"不能乱序)。

Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"partitioner", "murmur2"},  // ★ 哈希算法: murmur2 (和 Java Kafka 客户端默认一致)
});

// 用 user_id 做 key → 同一用户的订单永远进同一分区
producer.produce(MessageBuilder("orders")
    .key("user-1001")     // ★ 有 key → hash("user-1001") % partitionCount
    .payload(order_data)
);

什么是 murmur2? 一种非加密哈希函数,速度快、分布均匀。Kafka 用它来计算 key 的哈希值。它不是加密用的(不需要抗碰撞),只需要把 key 均匀地打散到各分区。

⚠️ 分区数变化会破坏顺序:如果分区数从 3 变成 5,hash("user-1001") % 5 的结果可能和 % 3 不一样,同一 key 的消息就会进不同分区。所以分区数只能加不能减,加了之后也要注意顺序问题。

策略三:无 Key 分区

消息不带 key 时,Kafka 需要自己决定进哪个分区。有两种方式:

Round-Robin(轮询):按顺序轮流来——第 1 条进分区 0,第 2 条进分区 1,第 3 条进分区 2,第 4 条又进分区 0……好处是均匀分布,坏处是每条消息可能进不同分区,没有顺序保证。

Sticky 分区器(Kafka 2.4+ 默认):不是严格轮询,而是"粘"在一个分区上一段时间——先把消息都往分区 0 发,攒满一个批量后再换到分区 1。为什么?因为如果严格轮询,每条消息进不同分区,每个分区都攒不满一个批量,就得发很多小批量。粘在同一个分区上,批量更容易攒满,网络请求更少,吞吐更高。

// 无 key → Sticky 分区(默认),均匀分布到各分区
producer.produce(MessageBuilder("logs")
    .payload(log_data)    // ★ 没有 key
);

怎么选? 需要顺序保证 → 用 key。不需要顺序、追求均匀 → 不带 key(Sticky 分区器自动处理)。需要精确控制 → 直接指定分区(少用)。

7.2 acks 机制

Producer 发消息后,要不要等 Broker 确认"收到了"?等多久?这就是 acks 机制。

为什么需要 acks? 网络不可靠——你发了消息,不代表对方收到了。如果不等确认就认为发送成功,消息可能在路上丢了你都不知道。但等确认又要花时间,所以要在可靠性和性能之间做取舍。

三种 acks 级别:

acks=0:发完就不管了,不等任何确认。最快,但 Broker 宕机或网络断了消息就丢了。适合日志收集这种"丢几条无所谓"的场景。

acks=1:等 Leader 写入成功后回一个确认。但如果 Leader 刚确认就挂了,Follower 还没来得及同步这条消息,消息就丢了。适合可以容忍少量丢失的场景。

acks=all(-1):等 Leader 和所有ISR副本都写入成功后才回确认。最安全,但最慢。生产环境推荐。

什么是 ISR? ISR(In-Sync Replicas)= "和 Leader 保持同步的副本集合"。一个分区有 1 个 Leader + N 个 Follower,但不是所有 Follower 都跟得上——如果某个 Follower 网络慢了或挂了,它会被踢出 ISR。只有 ISR 里的副本才算"可靠的"。

acks=all 的陷阱:如果 ISR 只剩 Leader 自己(其他副本都掉队了),acks=all 实际等价于 acks=1——Leader 一个人确认就行。这时 Leader 挂了还是会丢数据。所以需要配合 min.insync.replicas 使用——比如设为 2,则 ISR 里的副本数 < 2 时直接拒绝写入,宁可不可用也不丢数据。

// 生产环境推荐: acks=all + min.insync.replicas=2
Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"acks", "all"},           // ★ 等 ISR 所有副本确认
    {"retries", 2147483647},   // ★ 无限重试 (配合 delivery.timeout.ms)
    {"delivery.timeout.ms", 120000},  // ★ 总投递超时 2 分钟
    {"enable.idempotence", true},     // ★ 幂等生产者 (见 7.4)
    {"max.in.flight.requests.per.connection", 5},  // ★ 幂等模式下 ≤5
});

7.3 批量发送与压缩

Kafka 生产者默认批量发送——消息先进入RecordAccumulator(缓冲区),按分区攒成批量后一次性发送:

Producer 消息流转: 消息1 ─┐ 消息2 ─┤→ RecordAccumulator (按分区分组) 消息3 ─┤ ↓ 消息4 ─┘ 批量满 (batch.size=16KB) 或 超时 (linger.ms=0) ↓ Sender 线程发送到 Broker
参数 默认值 说明
batch.size 16384 (16KB) 每个批量的最大大小
linger.ms 0 攒批等待时间。0=立即发,>0=等这么久凑批量
buffer.memory 33554432 (32MB) 生产者总缓冲区大小
compression.type none 压缩算法: none/gzip/snappy/lz4/zstd
// 调优示例: 高吞吐场景
Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"batch.size", 65536},          // ★ 64KB 批量
    {"linger.ms", 10},              // ★ 最多等 10ms 凑批量
    {"compression.type", "lz4"},    // ★ LZ4 压缩 (速度最快)
    {"buffer.memory", 67108864},    // ★ 64MB 缓冲区
});

压缩的收益:压缩在生产者端做,解压在消费者端做,Broker 存储的是压缩后的数据。好处是网络带宽 + 存储空间 + 页缓存利用率都提升。LZ4 是速度和压缩比的最佳折中,生产环境首选。

7.4 幂等生产者(Idempotent Producer)

问题:Producer 发消息后网络超时,重试发送了同一条消息 → Broker 收到两份 → 消费者消费两次。

解法:开启幂等生产者,Broker 自动去重。

Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"enable.idempotence", true},   // ★ 开启幂等
    {"acks", "all"},                // ★ 幂等要求 acks=all
    {"retries", 2147483647},        // ★ 无限重试
    {"max.in.flight.requests.per.connection", 5},  // ★ ≤5 (幂等限制)
});

原理:

1. Producer 启动时获得一个 PID (Producer ID) 2. 每条消息携带: <PID, Partition, SequenceNumber> 3. Broker 收到后检查: - 如果 SequenceNumber = 上次的 +1 → 正常, 接受 - 如果 SequenceNumber ≤ 上次的 → 重复, 丢弃但返回成功 - 如果 SequenceNumber > 上次的 +1 → 乱序, 拒绝

限制:幂等只保证单分区单会话内不重复。Producer 重启后 PID 变化,无法跨会话去重。跨会话需要事务。

7.5 事务(Transactional Producer)

事务保证跨分区、跨会话的 Exactly Once 语义:

// C++ (librdkafka 事务 API)
Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"transactional.id", "order-tx-1"},  // ★ 事务 ID (跨会话不变)
    {"enable.idempotence", true},        // ★ 事务要求幂等
});

producer.init_transactions();  // ★ 初始化事务

while (true) {
    producer.begin_transaction();  // ★ 开启事务

    // 发送多条消息到不同分区/Topic
    producer.produce(MessageBuilder("orders").partition(0).payload(order_data));
    producer.produce(MessageBuilder("orders").partition(1).payload(order_data2));
    producer.produce(MessageBuilder("inventory").payload(stock_data));

    // 发送消费位移提交 (消费-处理-生产 模式)
    producer.send_offsets_to_transaction(consumer_offsets, consumer_group_id);

    // 提交或中止
    if (success) {
        producer.commit_transaction();    // ★ 全部原子提交
    } else {
        producer.abort_transaction();     // ★ 全部回滚
    }
}

事务的核心:transactional.id 是跨会话的稳定标识。Producer 重启后,Broker 会发现同一个 transactional.id 的上一个事务未完成,执行事务恢复(abort 或 commit pending 事务),然后才允许新 Producer 工作。这保证了跨会话的 Exactly Once。

消费端:消费者要读事务消息,需设置 isolation.level=read_committed,只读取已提交的事务消息。


8 · 消费者深入

8.1 消费者拉取模型

Kafka 消费者是 Pull 模型——消费者主动向 Broker 发起 Fetch 请求:

Consumer Broker │ │ │── FetchRequest(offset=N) ────→│ │ │ (查 Log Segment) │←── FetchResponse(msgs) ───────│ │ │ │ 处理消息... │ │ │ │── FetchRequest(offset=N+M) ──→│
对比 Push(Broker 推) Pull(Consumer 拉)
速率控制 Broker 控制,可能压垮 Consumer Consumer 自己控制
空轮询 无 有(没消息时白请求)
批量 Broker 决定 Consumer 决定

空轮询的解法:fetch.min.bytes + fetch.max.wait.ms。Consumer 告诉 Broker"至少给我 N 字节的数据,没有就等 M 毫秒再回"。这样没消息时 Broker 会 hold 住请求,有消息或超时才返回——长轮询。

8.2 Offset 管理

Kafka 0.9+ 的 Offset 存储在内部 Topic __consumer_offsets 中(不再用 ZooKeeper):

// 手动提交 (推荐)
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "order-processor"},
    {"enable.auto.commit", false},  // ★ 关闭自动提交
});

consumer.subscribe({"orders"});

while (true) {
    Message msg = consumer.poll();
    if (!msg || msg.get_error()) continue;

    json order = json::parse(msg.get_payload());
    try {
        process_order(order);
        consumer.commit(msg);  // ★ 同步提交: 阻塞等待 Broker 确认
    } catch (...) {
        // 不提交, 下次重新消费
    }
}
// 异步提交 (高吞吐场景)
consumer.commit_async(msg, [](const Error& err) {
    if (err) std::cerr << "提交失败: " << err << std::endl;
    // 异步提交失败不会重试, 可能导致重复消费
});

// 优雅退出时同步提交确保不丢
consumer.commit();  // 同步提交所有已消费但未提交的 offset
consumer.close();
提交方式 优点 缺点
自动提交 简单 可能重复消费、消息丢失(处理完还没提交就挂了,下次从旧 offset 开始)
同步提交 可靠 阻塞,降低吞吐
异步提交 不阻塞 提交失败可能重复消费

最佳实践:正常消费用异步提交(高吞吐),优雅退出时用同步提交(确保不丢)。处理完业务再提交 offset(At Least Once + 幂等)。

8.3 消费者位移重置

当消费者组第一次启动或 Offset 过期时,需要决定从哪里开始消费:

Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "new-group"},
    {"auto.offset.reset", "earliest"},  // ★ earliest=从头, latest=只读新的, none=报错
});
值 含义
earliest (-2) 从最早的可用消息开始
latest (-1) 只消费启动之后的新消息(默认)
none 没有已提交的 Offset 就报错

注意:auto.offset.reset 只在没有已提交 Offset 时生效。如果已经有提交过的 Offset,消费者会从上次的位置继续,不管这个配置。


9 · Rebalance 机制

图 4 \xb7 消费者组 Rebalance

9.1 什么时候触发 Rebalance

事件 说明
消费者加入 新消费者上线或重启
消费者离开 消费者主动关闭或被判定宕机(session.timeout.ms 超时)
订阅变化 消费者订阅的 Topic 列表变化
Topic 分区数变化 增加了分区

9.2 Rebalance 的代价

Rebalance 期间,所有消费者停止消费(Stop The World),直到分区重新分配完成。这是 Kafka 消费者最大的痛点:

  • 小集群:几百毫秒到几秒
  • 大集群(几百个消费者):可能几十秒到分钟级
  • 期间消息积压,延迟飙升

9.3 分配策略

策略 说明 优缺点
RangeAssignor 按 Topic 分,每个消费者分到连续的分区 可能不均匀
RoundRobinAssignor 跨 Topic 轮询分配 较均匀
StickyAssignor 尽量保持上次分配不变 ★ Rebalance 时迁移最少
CooperativeStickyAssignor 增量 Rebalance,不停所有消费者 ★★ 最优,Kafka 2.4+
// 推荐使用 CooperativeStickyAssignor
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "my-group"},
    {"partition.assignment.strategy", "cooperative-sticky"},  // ★ 增量 Rebalance
});

Cooperative Rebalance 的改进:传统 Rebalance 是"Eager 协议"——所有消费者先撤销全部分区,再重新分配。Cooperative 协议是"增量式"——只撤销需要迁移的分区,其他消费者继续消费。大幅减少 Stop The World 时间。

9.4 避免 Rebalance 的最佳实践

Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "my-group"},

    // ★ 心跳相关: 避免误判宕机
    {"session.timeout.ms", 30000},     // 心跳超时 30s (默认 10s 太短)
    {"heartbeat.interval.ms", 3000},   // 每 3s 发一次心跳 (建议 = session.timeout / 3)

    // ★ Poll 超时: 避免处理太慢被踢
    {"max.poll.interval.ms", 300000},  // 两次 poll 间隔最多 5 分钟 (默认 5 分钟)
    {"max.poll.records", 100},         // 每次 poll 最多 100 条 (减少处理时间)

    // ★ 分配策略
    {"partition.assignment.strategy", "cooperative-sticky"},
});

最常见的 Rebalance 原因:不是消费者真的挂了,而是 max.poll.interval.ms 超时——消费者处理消息太慢,两次 poll() 间隔超过了这个时间,Broker 以为消费者死了,触发 Rebalance。解法:调大 max.poll.interval.ms,或减少 max.poll.records 让每次处理少一点。


10 · 消费者组与分区分配详解

10.1 分配过程

1. 消费者加入 Group → 向 Group Coordinator 发 JoinGroup 请求 2. Coordinator 选一个消费者做 Group Leader 3. Group Leader 收到所有成员信息 → 计算分区分配方案 4. Leader 把分配方案发给 Coordinator 5. Coordinator 把分配结果下发给所有成员 6. 每个消费者收到自己的分区分配 → 开始消费

10.2 Group Coordinator

每个消费者组有一个 Group Coordinator——由某个 Broker 担任,负责:

  • 管理消费者组成员(加入/离开/心跳)
  • 触发和管理 Rebalance
  • 管理 Offset 提交
# 查看消费者组详情
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group order-processor

# 输出:
# TOPIC    PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID
# orders   0          12345           12400           55   consumer-1-xxx
# orders   1          6789            6800            11   consumer-2-xxx
# orders   2          9999            10000           1    consumer-1-xxx

Lag = LOG-END-OFFSET - CURRENT-OFFSET,表示积压了多少消息没消费。这是 Kafka 监控的核心指标。


第四部分 · 副本与高可用

这部分讲 Kafka 的副本机制、Leader 选举、高水位。读完你能搞懂 Kafka 怎么做到高可用和数据一致性。


11 · 副本机制与 ISR

图 2 \xb7 副本与 ISR 机制

11.1 副本模型

每个 Partition 有 1 个 Leader + N 个 Follower:

  • Leader:处理所有读写请求(生产者写、消费者读都走 Leader)
  • Follower:从 Leader 拉取数据同步,不处理客户端请求
Partition 0 (replication-factor=3): Broker 1: Leader ←── 所有读写 Broker 2: Follower ←── 从 Leader 拉取同步 Broker 3: Follower ←── 从 Leader 拉取同步

为什么 Follower 是拉而不是推? 拉模型让 Follower 自己控制同步速率,不会因为 Leader 推太快而压垮 Follower。而且 Follower 可以批量拉取,效率更高。

11.2 ISR(In-Sync Replicas)

ISR 是与 Leader 保持同步的副本集合,是 Kafka 高可用的核心:

集合 含义
AR (Assigned Replicas) 分区的所有副本
ISR (In-Sync Replicas) 与 Leader 同步的副本子集(包含 Leader 自己)
OSR (Out-of-Sync Replicas) 落后太多的副本
AR = ISR + OSR

进入 ISR 的条件:Follower 在 replica.lag.time.max.ms(默认 30 秒)内追上 Leader 的最新 Offset。

被踢出 ISR 的条件:Follower 超过这个时间没追上 Leader。

# 查看 Topic 的 ISR 状态
kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic orders

# 输出:
# Topic: orders  Partition: 0  Leader: 1  Replicas: 1,2,3  Isr: 1,2,3
#                                                          ↑ ISR 列表
# 如果 Broker 3 的副本掉队:
# Topic: orders  Partition: 0  Leader: 1  Replicas: 1,2,3  Isr: 1,2
#                                                          ↑ 3 被踢出

11.3 min.insync.replicas

# 创建 Topic 时指定最小同步副本数
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic orders \
  --partitions 3 --replication-factor 3 \
  --config min.insync.replicas=2  # ★ ISR 至少 2 个才允许写入

配合 acks=all 使用:acks=all + min.insync.replicas=2 + replication-factor=3 = 挂 1 个副本仍可写、挂 2 个副本拒绝写入。这是生产环境的黄金配置——在可用性和可靠性之间取得平衡。


12 · Leader 选举与高水位

12.1 Leader 选举流程

当 Leader 所在的 Broker 宕机时,Controller 会从 ISR 中选一个 Follower 成为新 Leader:

1. Broker 1 (Leader of P0) 宕机 2. ZooKeeper/KRaft 检测到 Broker 1 心跳超时 3. Controller 收到 Broker 1 下线通知 4. Controller 查找 P0 的 ISR: [1, 2, 3] → 去掉 1 → [2, 3] 5. 选 ISR 中的第一个 (Broker 2) 作为新 Leader 6. Controller 更新元数据, 通知所有 Broker 7. 生产者/消费者自动切换到新 Leader

12.2 Unclean Leader Election

问题:如果 ISR 中所有副本都挂了,怎么办?

配置 行为
unclean.leader.election.enable=false(默认) 分区不可用,等 ISR 中的副本恢复
unclean.leader.election.enable=true 从 OSR 中选一个非同步副本做 Leader,丢数据

生产环境必须设为 false。宁可暂时不可用也不能丢数据。可用性和数据一致性之间,Kafka 默认选择一致性。

12.3 高水位(High Watermark, HW)

高水位是 Kafka 保证数据一致性的核心机制:

Partition 0 的 Offset 轴: 0 ─────── 100 ─────── 150 ─────── 200 ─────── 250 ──→ │ │ │ │ │ │ LEO(Leader) │ LEO(Follower1) │ │ │ │ │ HW ────┘ │ │ │ │ │ 消费者只能读到 HW 以下 │ │ (HW = min(LEO of all ISR)) │
概念 含义
LEO (Log End Offset) 每个副本的日志末尾 Offset(下一条要写的位置)
HW (High Watermark) 所有 ISR 副本中最小的 LEO,消费者只能读 HW 以下的消息
HW 的作用 保证消费者只读到"已同步到所有 ISR 副本"的消息

为什么需要 HW? 如果消费者能读到 Leader 的 LEO,而某个 Follower 还没同步到那里,此时 Leader 挂了、Follower 成为新 Leader,消费者读到的消息就"消失"了。HW 保证了消费者读到的消息一定在所有 ISR 副本中都存在。

12.4 Leader Epoch(解决 HW 的 corner case)

HW 机制有一个已知问题:Leader 切换 + Follower 恢复时可能导致数据不一致或丢数据。Kafka 引入了 Leader Epoch 机制来解决:

leader-epoch-checkpoint 文件: 0 0 ← Epoch 0, 起始 Offset 0 1 150 ← Epoch 1, 起始 Offset 150 (第1次 Leader 切换) 2 200 ← Epoch 2, 起始 Offset 200 (第2次 Leader 切换)
  • 每次 Leader 切换,Epoch 号 +1
  • Follower 恢复时,先发 OffsetForLeaderEpoch 请求给 Leader,确认正确的截断点
  • 替代了单纯基于 HW 的截断逻辑,避免了数据丢失

一句话:HW 是"消费者能读到哪里"的保证,Leader Epoch 是"Leader 切换后副本怎么对齐"的保证。两者配合,Kafka 在大多数场景下做到了不丢数据。


13 · 高可用最佳实践

13.1 推荐配置

# Topic 级别
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic critical-events \
  --partitions 6 --replication-factor 3 \
  --config min.insync.replicas=2 \
  --config unclean.leader.election.enable=false \
  --config retention.ms=259200000  # 保留 3 天

# Producer 级别
# acks=all + enable.idempotence=true + retries=MAX

# Broker 级别
# default.replication.factor=3
# min.insync.replicas=2
# unclean.leader.election.enable=false

13.2 容灾能力

配置 挂 1 个 Broker 挂 2 个 Broker 挂 3 个 Broker
RF=1, minISR=1 丢数据 丢数据 丢数据
RF=2, minISR=1 可用,可能丢 不可用 不可用
RF=2, minISR=2 不可用 不可用 不可用
RF=3, minISR=2 可用,不丢 不可用,不丢 不可用
RF=3, minISR=1 可用,可能丢 可用,可能丢 不可用

RF=3, minISR=2 是业界标准:容忍 1 个节点故障且不丢数据。如果要容忍 2 个节点故障,需要 RF=5, minISR=3。


第五部分 · 性能调优与运维

这部分讲生产环境的调优、监控、运维。读完你能把 Kafka 跑到最佳状态,并在出问题时快速定位。


14 · 性能调优

14.1 生产者调优

场景 关键参数 说明
高吞吐 batch.size=65536, linger.ms=10, compression.type=lz4 攒大批量 + 压缩
低延迟 linger.ms=0, batch.size=16384, acks=1 不攒批,快速发送
高可靠 acks=all, enable.idempotence=true, retries=MAX 不丢不重
大消息 max.request.size=10485760 (10MB) 默认只支持 1MB
// 高吞吐配置模板
Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"acks", "all"},
    {"batch.size", 65536},
    {"linger.ms", 10},
    {"compression.type", "lz4"},
    {"buffer.memory", 67108864},
    {"enable.idempotence", true},
    {"max.in.flight.requests.per.connection", 5},
});

// 低延迟配置模板
Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"acks", "1"},
    {"linger.ms", "0"},
    {"batch.size", 16384},
    {"compression.type", "none"},
    {"enable.idempotence", true},
});

14.2 Broker 调优

# server.properties 关键调优

# 日志相关
log.dirs=/data/kafka-logs          # ★ 用独立磁盘 (SSD/NVMe)
num.io.threads=8                   # IO 线程数 (≈磁盘数)
num.network.threads=3              # 网络线程数 (默认3, 高负载调到8+)
socket.send.buffer.bytes=102400    # Socket 发送缓冲区
socket.receive.buffer.bytes=102400 # Socket 接收缓冲区

# 副本相关
num.replica.fetchers=4             # 副本拉取线程数 (默认1, 调大加速同步)
replica.fetch.max.bytes=1048576    # 副本拉取最大字节数

# 日志清理
log.retention.hours=168            # 保留 7 天
log.segment.bytes=1073741824       # Segment 大小 1GB
log.cleanup.policy=delete          # 删除策略

# 页缓存
# ★ 不要给 JVM 分配太多内存! Kafka 依赖页缓存
# 推荐: JVM 堆 6~8GB, 剩余给 OS 页缓存
# KAFKA_HEAP_OPTS="-Xmx6g -Xms6g"

最重要的 Broker 调优:不要给 Kafka JVM 分配太多堆内存。Kafka 不靠 JVM 堆存数据,靠 OS 页缓存。一台 32GB 内存的机器,JVM 给 6GB,剩下 26GB 全给页缓存——这才是正确的做法。

14.3 消费者调优

// 高吞吐消费者
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "my-group"},
    {"enable.auto.commit", false},
    {"fetch.min.bytes", 10240},         // ★ 至少拉 10KB
    {"fetch.max.wait.ms", 500},         // ★ 最多等 500ms
    {"max.poll.records", 500},          // ★ 每次 poll 500 条
    {"max.partition.fetch.bytes", 1048576},  // ★ 每分区最多 1MB
});

14.4 操作系统层面

# 1. 文件描述符
ulimit -n 100000
# /etc/security/limits.conf:
# kafka  soft  nofile  100000
# kafka  hard  nofile  100000

# 2. Socket 缓冲区
sysctl -w net.core.rmem_default=262144
sysctl -w net.core.wmem_default=262144
sysctl -w net.core.rmem_max=16777216
sysctl -w net.core.wmem_max=16777216

# 3. 减少 swap
sysctl -w vm.swappiness=1
# Kafka 不应该被 swap 到磁盘!

# 4. 页缓存刷脏页策略
sysctl -w vm.dirty_ratio=80
sysctl -w vm.dirty_background_ratio=5
# 让 OS 多攒脏页再刷, 提高写入吞吐

# 5. 文件系统: XFS > ext4 (Kafka 官方推荐 XFS)

15 · 监控指标

15.1 核心监控指标

指标 来源 告警阈值 说明
Under-Replicated Partitions Broker JMX > 0 副本不同步的分区数,最重要的指标
Offline Partitions Broker JMX > 0 没有可用 Leader 的分区数,严重告警
Active Controller Count Broker JMX ≠ 1 Controller 数量,必须恰好 1 个
Consumer Group Lag Consumer > 10000 消费积压
ISR Shrink/Expand Broker JMX 频繁变化 ISR 频繁变动说明网络/磁盘有问题
Request Latency Broker JMX > 100ms 请求处理延迟
Disk Usage OS > 80% 磁盘空间
Network Usage OS > 80% 网络带宽
# 快速检查集群健康
kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-replicated-partitions
# 有输出 = 有副本不同步, 需要排查

kafka-topics.sh --bootstrap-server localhost:9092 --describe --unavailable-partitions
# 有输出 = 有分区不可用, 严重!

15.2 JMX 监控

# 启动 Kafka 时开启 JMX
export JMX_PORT=9999
export KAFKA_JMX_OPTS="-Dcom.sun.management.jmxremote=true \
  -Dcom.sun.management.jmxremote.port=9999 \
  -Dcom.sun.management.jmxremote.authenticate=false \
  -Dcom.sun.management.jmxremote.ssl=false"

生产环境推荐:用 Prometheus + JMX Exporter + Grafana 搭建监控。社区有现成的 Kafka Grafana Dashboard 模板。


16 · 常见运维操作

16.1 分区重分配(扩容)

新增 Broker 后,需要把旧 Broker 上的分区迁移到新 Broker:

# 1. 生成迁移计划
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
  --topics-to-move-json-file topics.json \
  --broker-list "4,5" \  # 新 Broker ID
  --generate

# 2. 执行迁移
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
  --reassignment-json-file plan.json \
  --execute

# 3. 验证迁移状态
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
  --reassignment-json-file plan.json \
  --verify

注意:分区迁移会消耗大量网络和磁盘 I/O,建议在低峰期执行。可以配合 kafka-reassign-partitions.sh --throttle 限速。

16.2 Preferred Leader 选举

副本迁移后,Leader 可能不在"首选副本"上。执行 Preferred Leader 选举让 Leader 回到首选位置:

kafka-leader-election.sh --bootstrap-server localhost:9092 \
  --election-type PREFERRED \
  --all-topic-partitions

16.3 滚动升级

1. 检查集群健康 (no under-replicated partitions) 2. 逐个 Broker: a. 优雅停止 (等待 Follower 同步完成) b. 更新配置/版本 c. 重启 d. 等待 ISR 恢复 (回到正常状态) 3. 全部完成后验证
# 优雅停止 Broker (等待副本同步)
kafka-server-stop.sh
# Kafka 会等待 controlled.shutdown=true 的 Broker 同步完再停

# 检查是否还有 under-replicated partitions
kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-replicated-partitions
# 确认输出为空后再继续下一个 Broker

17 · 常见故障与排查

17.1 消息积压(Lag 飙升)

排查路径: Lag 高 → 消费者处理慢? → 加消费者 (但 ≤ 分区数) → 优化消费逻辑 (批量写 DB、异步处理) → 临时扩分区 (但会影响顺序) → 降级非核心消费 (跳过部分消息)

17.2 副本不同步(Under-Replicated)

排查路径: ISR 缩小 → Follower 追不上 Leader? → 网络问题? (检查 Broker 间网络延迟) → 磁盘慢? (检查 iostat, Follower 磁盘 IO) → GC 停顿? (检查 JVM GC 日志) → 副本拉取线程不够? (调大 num.replica.fetchers)

17.3 Producer 发送超时

排查路径: 发送超时 → Broker 响应慢? → Broker 负载高? (检查 CPU/IO/网络) → 磁盘满? (检查 df -h) → 网络问题? (检查 Producer 到 Broker 的网络) → 批量太大? (减小 batch.size) → 缓冲区满? (增大 buffer.memory)

17.4 Rebalance 频繁

排查路径: 频繁 Rebalance → 消费者频繁进出? → session.timeout.ms 太短? (调大) → max.poll.interval.ms 太短? (调大) → 消费者处理太慢? (减少 max.poll.records) → 消费者 OOM? (检查内存) → 用 CooperativeStickyAssignor (减少 STW)

第六部分 · Kafka 生态


18 · Kafka Streams

Kafka Streams 是 Kafka 内置的流处理库——不需要独立的流处理集群(如 Flink/Spark),直接在应用内处理:

// Java 示例 (Kafka Streams 是 Java 库, 没有 C++ 版)
StreamsBuilder builder = new StreamsBuilder();

// 从 orders Topic 读取
KStream<String, String> orders = builder.stream("orders");

// 按用户分组, 统计订单数
KTable<String, Long> userOrderCounts = orders
    .groupByKey()
    .count();

// 写回 Kafka
userOrderCounts.toStream().to("user-order-counts");

StreamsConfig config = new StreamsConfig(properties);
KafkaStreams streams = new KafkaStreams(builder.build(), config);
streams.start();
特性 Kafka Streams Flink
部署 嵌入应用,无独立集群 独立集群
依赖 只依赖 Kafka 需要 Flink JobManager/TaskManager
状态管理 RocksDB 本地状态 RocksDB
Exactly Once Kafka 事务 Checkpoint + 两阶段提交
适合 Kafka 内部流处理 复杂流处理、窗口计算

选型:如果你的数据已经在 Kafka 里,处理逻辑不复杂(过滤、聚合、join),用 Kafka Streams 最简单——不需要额外集群。如果需要复杂窗口、CEP、批流统一,用 Flink。


19 · Kafka Connect

Kafka Connect 是 Kafka 的数据导入/导出框架——标准化地连接 Kafka 和外部系统:

Source Connector (导入): Sink Connector (导出): MySQL ──→ Kafka Kafka ──→ Elasticsearch MongoDB ──→ Kafka Kafka ──→ S3 File ──→ Kafka Kafka ──→ JDBC
# 启动 Connect (单机模式)
connect-standalone.sh config/connect-standalone.properties \
  connectors/mysql-source.properties \
  connectors/es-sink.properties

Connect 的价值:不用自己写 Producer/Consumer 代码来搬运数据。社区有几百个现成的 Connector(Debezium、JDBC、Elasticsearch、S3 等),配置即用。


20 · Schema Registry

Kafka 消息是字节流,没有内置 Schema。Schema Registry 解决"消息格式怎么管理"的问题:

Producer → [序列化 + 注册 Schema] → Kafka (存 Schema ID + Avro/Protobuf 数据) Consumer → [读 Schema ID → 从 Registry 获取 Schema → 反序列化]
能力 说明
Schema 存储 集中管理 Avro/JSON/Protobuf Schema
版本管理 Schema 变更有版本号,可追溯
兼容性检查 新 Schema 必须与旧版本兼容(向后/向前/完全兼容)
减小消息体积 消息只存 Schema ID(4 字节),不存 Schema 本身

为什么需要 Schema Registry? 没有它,Producer 改了消息格式(加字段/改字段类型),Consumer 不知道,反序列化直接崩。Schema Registry 在 Producer 发送前检查兼容性,不兼容直接拒绝,从源头防止格式不一致。


21 · 总结

21.1 Kafka 的核心设计哲学

设计 哲学
日志即数据库 消息不是"取走就没了",而是持久化的、可回放的事件日志
分区即并行 分区是并行度的最小单位,加分区 = 加并行
拉模型 消费者主动拉取,自己控制速率
页缓存优先 不靠 JVM 堆,靠 OS 页缓存
简单即快 顺序写 + 稀疏索引 + 零拷贝,没有复杂的数据结构

21.2 什么时候该用 Kafka

该用:

  • 日志收集、事件溯源、流处理
  • 高吞吐场景(万级到百万级 TPS)
  • 需要消息回放(重新消费历史消息)
  • 多消费者独立消费同一份数据(不同 Consumer Group)
  • 数据管道(Kafka Connect + Streams)

不该用:

  • 需要复杂路由(RabbitMQ 更合适)
  • 需要事务消息(RocketMQ 更合适)
  • 消息量小、团队没有运维能力
  • 需要严格的全局有序(Kafka 只能分区有序)
  • 需要消息优先级(Kafka 不支持)

21.3 学习建议

  1. 先跑通——Docker 起一个 KRaft 模式 Kafka,写 Producer + Consumer 收发消息
  2. 再理解架构——搞懂 Partition、Replica、ISR、Controller 的关系
  3. 然后深入存储——看 Partition 目录下的文件,理解 Log Segment 和索引
  4. 最后看运维——故意 kill 掉 Broker,观察 Leader 选举和 Rebalance
  5. 生产加固——配置 acks=all + min.insync.replicas=2,监控 Lag 和 ISR

现实节奏:Kafka 入门需要 1~2 周(基本收发 + 核心概念),深入需要 1~2 个月(存储引擎 + 副本机制 + 调优),精通需要半年以上(生产运维 + 故障排查 + 架构设计)。不要急,先把基本概念搞透,再在实践中逐步深入。

本页目录