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 · 消息全生命周期
在深入每个组件之前,先跟着一条消息走一遍完整旅程。这一节只求建立直觉,不展开细节——每个步骤后面都有专门章节深入讲解。

生产端:消息怎么发出去的
假设你是一个 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.1 核心组件
| 组件 |
职责 |
类比 |
| Broker |
存储和转发消息的服务节点 |
邮局的分拣中心 |
| Producer |
发送消息的客户端 |
寄信人 |
| Consumer |
消费消息的客户端 |
收信人 |
| Controller |
管理集群状态(分区/副本/Leader 选举) |
邮局局长 |
| ZooKeeper / KRaft |
元数据存储 + Controller 选举 |
邮局档案室 |
1.2 一条消息的生命周期

关键点:
- 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:

| 对比 |
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 · 索引文件与零拷贝

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
查找流程:
- 二分查找
.index 文件,定位到 ≤ 目标 offset 的最近索引项
- 跳到
.log 文件对应的物理位置
- 顺序扫描
.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 机制

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

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 学习建议
- 先跑通——Docker 起一个 KRaft 模式 Kafka,写 Producer + Consumer 收发消息
- 再理解架构——搞懂 Partition、Replica、ISR、Controller 的关系
- 然后深入存储——看 Partition 目录下的文件,理解 Log Segment 和索引
- 最后看运维——故意 kill 掉 Broker,观察 Leader 选举和 Rebalance
- 生产加固——配置 acks=all + min.insync.replicas=2,监控 Lag 和 ISR
现实节奏:Kafka 入门需要 1~2 周(基本收发 + 核心概念),深入需要 1~2 个月(存储引擎 + 副本机制 + 调优),精通需要半年以上(生产运维 + 故障排查 + 架构设计)。不要急,先把基本概念搞透,再在实践中逐步深入。