13794 字
69 分钟
Kafka Consumer Group:从 poll()、Offset 到 Rebalance 的分布式协调系统

Kafka Consumer Group:从 poll()、Offset 到 Rebalance 的分布式协调系统#

深入 Fetch Runtime、Group Coordinator、Partition Ownership、Offset Commit 与 Kafka 4.x 新 Consumer Protocol

阅读目标#

前三篇已经依次建立了三层模型:

第一篇
Log → Partition → Offset → Consumer Position
第二篇
RecordBatch → LogSegment → Index → Broker Data Path
第三篇
KafkaProducer.send()
→ RecordAccumulator
→ Sender
→ NetworkClient
→ Partition Leader

第四篇把视角移动到消费端。

很多 Kafka 教程会把 Consumer 简化成:

while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
handle(record);
}
}

这段代码看起来只是一个循环,但 poll() 背后同时推进了多套状态:

Partition 数据获取
Consumer 本地缓存
反序列化与交付
Group Coordinator 查找
成员身份与存活
Partition Assignment
Rebalance Callback
Offset Fetch / Commit
Metadata 与网络 I/O

因此 Consumer 最难的部分不是“怎么拉消息”,而是:

Kafka 如何在一个动态变化的分布式成员集合中,持续维护 Partition 的唯一 Ownership,并把业务处理进度安全地映射为可恢复的 Offset。

完成本篇后,你应该能够:

  • 解释为什么 Kafka Consumer 使用 Pull,而不是由 Broker 主动 Push;
  • 说明一次 poll() 为什么不等于一次 FetchRequest;
  • 画出 Consumer 的 Fetch、缓存、poll 与业务处理流水线;
  • 区分 Record Offset、Consumer Position、Committed Offset、LSO、HW 与 LEO;
  • 解释 subscribe()assign() 的所有权模型差异;
  • 解释同一 Consumer Group 内 Partition 与 Consumer 的分配约束;
  • 从 Group ID 推导 Group Coordinator 与 __consumer_offsets 的职责;
  • 画出 Classic Protocol 的 FindCoordinator → JoinGroup → SyncGroup 完整链;
  • 解释 Range、RoundRobin、Sticky 与 CooperativeSticky 的不同权衡;
  • 解释 Eager Rebalance 为什么会产生全组停顿;
  • 说明 Kafka 4.x 新 Consumer Protocol 为什么不再依赖全局同步屏障;
  • 区分 heartbeat.interval.mssession.timeout.msmax.poll.interval.ms
  • 解释为什么 Heartbeat 正常仍可能发生 Livelock;
  • 正确设计手动 Offset Commit,避免提交过早或提交乱序;
  • 根据 Consumer Metrics 和 Group 状态定位 Lag 与 Rebalance Storm;
  • 沿 Classic Consumer 与 Async Consumer 两条源码路线继续深入。

本文以 Apache Kafka 4.3.x 官方文档与 Java Client API 为主线。Kafka 4.0 起,KIP-848 新 Consumer Rebalance Protocol 已经 GA,但客户端默认仍是 group.protocol=classic;使用新协议需要显式设置 group.protocol=consumer。新协议采用 Broker 驱动的分配和 Fully Incremental Reconciliation,并使用改进后的 Consumer 线程模型。1


0. poll() 返回空集合,真的只是没有消息吗#

先看最普通的调用:

ConsumerRecords<String, String> records =
consumer.poll(Duration.ofSeconds(1));

如果返回:

records.isEmpty() == true

初学者最容易得出:

Topic 里现在没有消息。

但空集合可能来自完全不同的状态:

本地 Completed Fetch 暂时为空
FetchRequest 已发出但尚未返回
当前 Consumer 还没有获得 Partition Assignment
Group 正在 Rebalance
目标 Position 已经追到可见上界
read_committed 被 LSO 阻挡
Partition 被 pause
Metadata 或 Leader 暂时不可用

poll 返回空集合并不等于 Topic 没有消息

更重要的是,即便这一次没有返回 Record,poll() 仍可能在推进:

网络事件
Group Join
Assignment 回调
Offset 初始化
自动提交
Metadata 刷新
Fetch 调度

所以,理解 Consumer 的第一步是:

poll() 理解成 Consumer Runtime 的应用线程推进边界,而不是简单的一次“远程拉取消息”。

0.1 一个调用包含三类职责#

可以把 Consumer 工作拆成三条互相耦合但不同的链。

数据链#

FetchRequest
Broker Partition Leader
FetchResponse
Completed Fetch / Local Buffer
Deserialize
ConsumerRecords

成员链#

FindCoordinator
Join / Heartbeat
Member Identity
Session Liveness

Ownership 链#

Subscription
Assignment Calculation
Partition Revocation / Assignment
Reconciliation

Offset Commit 又把业务处理进度接回 Group Coordinator。

如果把这几条链混成一句:

Consumer 通过 poll 从 Broker 拉消息。

那么后面的 Rebalance、Lag、重复消费和超时问题几乎都会解释错。


1. Kafka Consumer 的对象与线程模型#

1.1 KafkaConsumer 不是线程安全的#

Kafka 官方 Javadoc 明确说明:

KafkaConsumer is NOT thread-safe

wakeup() 之外,不应该让多个线程无同步地调用同一个 Consumer。3

典型正确模型是:

一个 Consumer
一个 Poll Thread
顺序调用 Consumer API

关闭时,其他线程可以:

consumer.wakeup();

使正在阻塞的操作抛出 WakeupException,从而让 Poll Thread 进入退出逻辑。

不要把 Producer 的经验直接复制到 Consumer:

KafkaProducer
→ 线程安全,适合多个业务线程共享
KafkaConsumer
→ 非线程安全,Ownership、Position 与 Callback
都依赖单一应用线程顺序

1.2 为什么 Consumer 比 Producer 更依赖线程约束#

Consumer 内部维护大量与调用顺序相关的状态:

当前 Subscription
当前 Assignment
每个 Partition 的 Position
Pause 状态
Pending Commit
Rebalance Callback
本地 Fetch Buffer
Group Member / Epoch

例如:

consumer.seek(tp, 100);
consumer.pause(Set.of(tp));
consumer.commitSync();

这些操作不是无状态 RPC,而是在修改同一个消费状态机。

如果多个线程同时调用:

线程 A:poll()
线程 B:seek()
线程 C:commit()

应用很难定义可靠的先后语义。

1.3 Classic 与 Consumer 两种内部 Runtime#

Kafka 4.x 的同一个 KafkaConsumer API,可以运行在两种 Group Protocol 下。

Classic 与新 Consumer Protocol 的线程模型

Classic Protocol#

group.protocol=classic

概念主线:

Application Thread
ClassicKafkaConsumer
ConsumerCoordinator / SubscriptionState / Fetch
ConsumerNetworkClient
Broker
Heartbeat Thread
Group Coordinator

Classic Client 在本地承担较多协调职责:

  • 参与 JoinGroup / SyncGroup;
  • 上报 Assignor 能力;
  • Group Leader 计算 Assignment;
  • Heartbeat Thread 维持成员 Session;
  • 应用线程在 poll 中完成 Rebalance Callback 和状态推进。

Consumer Protocol#

group.protocol=consumer

概念主线:

Application Thread
AsyncKafkaConsumer
Application Event Queue
================ Thread Boundary ================
ConsumerNetworkThread
├── HeartbeatRequestManager
├── ConsumerMembershipManager
├── CommitRequestManager
├── FetchRequestManager
└── NetworkClientDelegate

新协议把更多网络与协调工作放入专门的 Consumer Network Thread,并使用 Application Event 与 Background Event 在应用线程和网络线程之间传递状态。

这并不意味着:

使用新协议以后,应用可以多线程调用 KafkaConsumer。

公开 API 的线程安全约束仍然存在。

1.4 线程模型改变了什么#

新线程模型的目标之一,是减少:

网络协调进度
对应用 poll 调用节奏的直接依赖

但业务处理仍必须尊重:

max.poll.interval.ms

因为 Kafka 不只关心“进程是否活着”,还要判断:

这个持有 Partition 的应用是否仍在推进消费循环。

这就是 Heartbeat 与 Poll Liveness 必须分开理解的原因。


2. poll() 的完整执行模型#

2.1 poll() 并不等于同步 Fetch#

错误模型:

poll()
发送 FetchRequest
等待 FetchResponse
返回这些 Record

更接近真实 Runtime 的模型:

检查 Subscription / Assignment
推进 Group Coordination
初始化缺失 Position
从本地 Completed Fetch 读取
必要时调度新的 FetchRequest
处理网络事件
反序列化可交付 Record
更新 Position
返回 ConsumerRecords

一次 poll() 可能:

  • 直接返回上一次 Fetch 已缓存的数据;
  • 发出新的 Fetch 但本轮没有数据可返回;
  • 在 Rebalance 后执行 Callback;
  • 从 Committed Offset 初始化 Position;
  • 因为 auto.offset.reset 进行 Offset Reset;
  • 只返回本地缓存中的一部分 Record。

2.2 Consumer 使用 Pull 的设计意义#

如果 Broker 主动 Push:

Broker 决定发送速度
Consumer 处理能力不足
客户端 Buffer 持续膨胀

Pull 模型使 Consumer 可以根据自己的节奏请求数据:

Consumer Capacity
Fetch Timing / Fetch Size

这天然支持:

  • Batch Fetch;
  • 不同 Consumer 的独立速度;
  • Consumer 侧 Backpressure;
  • 失败后按 Position 重读;
  • Broker 不需要维护每个消费者的业务处理窗口。

但 Pull 不代表完全没有预取。

Kafka Consumer 为了吞吐,会并行向多个 Broker 发 Fetch,并把结果缓存在客户端。

2.3 Fetch、缓存与 poll 的流水线#

poll、Fetch 与本地缓存流水线

典型流程:

FetchRequest N 已发出
应用正在处理 poll N-1 返回的 Record
FetchResponse N 到达本地缓存
下一次 poll 从缓存交付数据
同时调度 FetchRequest N+1

因此:

Network Fetch
Application Processing

可以形成 Pipeline。

2.4 max.poll.records 不限制底层 Fetch 大小#

Kafka 4.3 默认:

max.poll.records = 500

它控制:

一次 poll() 最多向应用返回多少条 Record。

但官方配置明确说明,它不会改变底层 Fetch 行为。Consumer 可以把 FetchResponse 缓存在本地,再通过多次 poll 分批返回。2

所以:

FetchResponse = 3,000 records
max.poll.records = 500

可能表现为:

poll #1 → 500
poll #2 → 500
...

而不是 Broker 必须发送六次请求。

这对慢消费调优非常重要:

  • 减少 max.poll.records 可以缩短单次业务处理时间;
  • 但不一定减少网络 Fetch 或客户端内存规模;
  • Fetch 内存仍受 fetch.max.bytesmax.partition.fetch.bytes 和 Batch 大小影响。

2.5 Fetch 的关键参数必须放回数据链#

fetch.min.bytes#

Broker 至少积累多少可返回数据后再响应。

更大:

吞吐和批量效率可能提升
延迟可能增加

fetch.max.wait.ms#

即使没有达到 fetch.min.bytes,Broker 最长等待多久。

二者共同形成:

满足最小字节
达到最长等待

max.partition.fetch.bytes#

每个 Partition 的目标返回上限。

但如果首个可读 RecordBatch 本身超过该值,Broker 仍可能返回完整 Batch,保证 Consumer 可以继续前进。2

fetch.max.bytes#

整个 FetchResponse 的目标总上限,同样不是不可突破的绝对单 Batch 上限。

2.6 Fetch 的并行单位#

Consumer 会根据当前 Assignment,把 Partition 按 Leader Broker 聚合:

P0 Leader → Broker 1
P1 Leader → Broker 2
P2 Leader → Broker 1

形成:

FetchRequest Broker 1
├── P0 offset=100
└── P2 offset=780
FetchRequest Broker 2
└── P1 offset=500

与 Producer 类似:

Partition 是顺序和 Position 单位
Broker Node 是网络请求聚合单位

2.7 poll() 返回后 Position 已经推进#

Kafka Javadoc 将 Position 定义为:

下一条将交付给应用的 Record Offset。

poll() 返回 Record 后,Consumer 的 Position 会自动前进。3

注意:

Position 已前进
业务处理已经成功
Committed Offset 已更新

这三个状态之间的间隔,是 At-Least-Once、At-Most-Once 和重复消费问题的根源。


3. Offset 坐标系:不要再把所有 Offset 混成一个数#

Consumer 工程中最常见的概念错误之一,是只会说:

当前 Offset 是 100。

专家必须先问:

你说的是哪个 Offset?

Consumer Offset 坐标系

3.1 Record Offset#

某条 Record 在特定 Partition Log 中的位置:

TopicPartition(order-events, 3)
Record Offset = 108

它的作用域永远是:

Topic + Partition

不同 Partition 可以同时存在 Offset=108。

3.2 Consumer Position#

Position 是:

下一条应该交付给当前 Consumer 的 Offset。

如果 poll 返回:

98, 99, 100

Position 通常推进到:

101

但 Offset 不保证绝对连续。

Log Compaction、事务控制记录或中止事务都可能让可见 Record 之间出现 Offset Gap。3

所以,不应该写:

nextOffset = lastOffset + 1;

然后假设下一条业务 Record 一定存在。

3.3 Committed Offset#

Committed Offset 是:

持久记录的恢复位置。

Consumer 重启或 Rebalance 获得 Partition 后,会从 Committed Offset 继续。

最重要的语义是:

Committed Offset
=
下一条应该重新读取的 Offset

如果已经成功处理 Offset 99,应该提交:

100

而不是 99。

3.4 Log End Offset#

Log End Offset(LEO)可概念化为:

Leader 本地 Log 的下一写入位置。

如果最后一条 Record Offset 为 149:

LEO = 150

LEO 包含 Leader 已追加但未必对普通 Consumer 可见的尾部范围。

3.5 High Watermark#

High Watermark(HW)表示:

普通 Consumer 可以安全读取的复制提交上界。

这里先理解为可见边界。

为什么 HW 不能简单等于 LEO,以及它如何跟随 Replica Fetch 推进,放到第五篇副本与可靠性中深入。

3.6 Last Stable Offset#

使用:

isolation.level=read_committed

时,Consumer 只返回已提交事务中的 Record。

如果 Log 中存在未完成事务,Consumer 可能只能读取到:

Last Stable Offset(LSO)

官方文档指出,read_committed Consumer 的可读终点是第一个 Open Transaction 之前的稳定位置,因此 LSO 可能落后于 HW。2

本篇只建立坐标。

Transaction Marker、Abort Index 与 LSO 推进机制在第五篇展开。

3.7 Lag 到底怎么算#

常见 Consumer Lag:

Lag
Log End / 可读上界
-
Committed Offset 或 Current Position

但不同监控工具可能使用不同坐标:

  • Broker Log End Offset;
  • Consumer 可见上界;
  • Current Position;
  • Committed Offset;
  • read_committed 下的 LSO。

因此看到:

Lag = 100,000

第一反应不应该是立刻扩容,而是确认:

这个 Lag 的右端点是什么?
左端点是 Committed 还是 Position?
是否存在事务阻塞?
是否存在暂停 Partition?
是否是单一 Hot Partition?

4. subscribe()assign():谁负责 Partition Ownership#

Kafka Consumer 有两种主要消费模式。

subscribe 与 assign 的职责边界

4.1 subscribe():Kafka 管理 Ownership#

consumer.subscribe(List.of("order-events"));

应用声明:

我希望消费这些 Topic

Kafka Consumer Group 负责决定:

这个实例具体拥有哪些 Partition

它会使用:

  • Group ID;
  • Group Coordinator;
  • Membership;
  • Assignment Strategy;
  • Rebalance;
  • Committed Offset。

适用于绝大多数需要:

自动扩容
故障接管
实例动态变化

的业务消费者。

4.2 assign():应用管理 Ownership#

consumer.assign(List.of(
new TopicPartition("order-events", 0),
new TopicPartition("order-events", 1)
));

此时动态分区分配和 Consumer Group Coordination 被禁用。3

应用自己负责:

  • 哪个实例处理哪个 Partition;
  • 实例数量变化;
  • Partition 扩容;
  • 故障接管;
  • Ownership 冲突;
  • Assignment 发布。

assign() 并不代表不能使用 Kafka Offset Commit。

如果配置 group.id,应用仍可提交 Offset,但 Kafka 不会替你动态分配 Partition。

4.3 两种模式不能混用#

一个 Consumer 不能同时:

subscribe()
+
assign()

除非先:

consumer.unsubscribe();

因为两种 API 的 Ownership 来源不同:

subscribe → Coordinator Assignment
assign → Application Assignment

4.4 什么时候适合 assign()#

典型场景:

  • 外部调度系统已经拥有 Shard Ownership;
  • 固定 Partition 与物理实例绑定;
  • 离线回放工具;
  • 特殊数据修复;
  • 测试和诊断;
  • 需要完全绕开 Group Rebalance 的独立处理器。

普通业务服务不应该为了“避免 Rebalance”随意改用 assign()

那只是把一个成熟协调协议替换成自行实现分布式 Ownership。


5. Consumer Group:Partition 如何分给一群消费者#

5.1 Group ID 定义消费视图#

group.id=order-fulfillment

所有使用同一 Group ID 的 Consumer 组成一个逻辑消费组。

同一 Topic 可以被多个 Group 独立消费:

Group fulfillment
Group risk
Group analytics
Group notification

每个 Group 拥有自己的:

Partition Assignment
Committed Offset
Lag
Processing Lifecycle

Consumer Group 的 Partition Ownership

5.2 同一 Group 内的核心约束#

对传统 Consumer Group:

一个 Partition 在任一时刻只分配给同一 Group 中的一个 Consumer。

例如:

Topic = 8 Partitions
Group = 3 Consumers

可能分配:

Consumer A → P0 P1 P2
Consumer B → P3 P4 P5
Consumer C → P6 P7

这保证:

  • Partition 内 Record 由单一 Consumer 顺序拉取;
  • Position 不会由同组多个实例竞争推进;
  • 故障时可以把 Ownership 转移给其他成员。

5.3 Consumer 数大于 Partition 数#

Partitions = 3
Consumers = 5

最多只有三个 Consumer 获得有效 Partition。

其余实例空闲。

因此:

有效并行 Consumer 数
<=
Partition 数

这也是为什么看到 Lag 高时,不能无脑增加 Consumer。

5.4 一个 Consumer 可以拥有多个 Partition#

Kafka 不要求一 Consumer 一 Partition。

Partitions = 12
Consumers = 3

每个 Consumer 大约拥有四个 Partition。

Consumer 会在多个 Partition 间 Fetch、缓存和返回 Record。

业务处理慢时,需要进一步判断:

  • 每个 Partition 都慢;
  • 只有一个 Hot Partition;
  • 一个 Consumer 拥有过多 Partition;
  • 外部依赖造成所有 Partition 阻塞。

5.5 Share Group 不是 Consumer Group 的简单新版#

Kafka 4.2 起 Share Groups 已生产可用,允许同一 Share Group 的多个消费者协作处理同一 Partition 中的 Record,并提供 Per-record Acknowledgement 与 Delivery Attempt Count。12

但它是另一套消费语义:

Consumer Group
→ Partition Ownership
Share Group
→ Record Acquisition / Acknowledgement

本专题仍以传统 Consumer Group 为主。

需要任务队列语义时,可以在扩展篇单独研究 Share Groups,不能把它和 KIP-848 Consumer Protocol 混为一谈。


6. Group Coordinator:谁维护成员和消费进度#

Consumer Group 是分布式成员集合,必须有一个协调点回答:

这个 Group 当前有哪些成员?
每个成员订阅什么 Topic?
当前 Assignment 是什么?
成员是否仍存活?
当前 Generation / Epoch 是什么?
每个 TopicPartition 的 Committed Offset 是多少?

这个角色就是:

Group Coordinator

6.1 Coordinator 是一个 Broker 角色#

Group Coordinator 并不是额外部署的独立服务。

某个 Broker 会成为特定 Consumer Group 的 Coordinator。

Consumer 首先通过:

FindCoordinator(group.id)

找到它。

之后,Join、Heartbeat、Commit 和 Offset Fetch 等 Group 请求都发往该 Broker。

6.2 Group ID 如何映射到 Coordinator#

简化理解:

hash(group.id)
__consumer_offsets Partition
该 Partition Leader Broker
Group Coordinator

因此 Coordinator 的分布与内部 Offset Topic Partition 相关。

大量 Group 可以分散到不同 Coordinator Broker。

6.3 Coordinator 与 Controller 不同#

不要混淆:

KRaft Controller#

管理集群 Metadata:

Broker
Topic
Partition
Replica
Leader

Group Coordinator#

管理消费组状态:

Member
Subscription
Assignment
Offset Commit
Group Epoch

一个属于集群控制面,一个属于消费组协调。

6.4 Coordinator 内存状态与持久状态#

Coordinator 会缓存:

Group Metadata
Committed Offset

以快速响应 Heartbeat、Join 与 OffsetFetch。

持久化的 Offset Commit 会写入:

__consumer_offsets

这是一个 Compacted Internal Topic。

当 Coordinator 发生迁移时,新 Coordinator 可以加载对应 Offset Topic Partition 重建缓存。4


7. Classic Protocol:FindCoordinator、JoinGroup 与 SyncGroup#

Classic Consumer Group 的核心不是“Kafka 自动分配”。

它是一套明确的协议。

Classic Consumer Group Rebalance 时序

7.1 FindCoordinator#

Consumer 先向任意 Broker 查询:

group.id 对应的 Coordinator 在哪里

获得:

Coordinator Node

后续 Group 请求发给该 Node。

如果 Coordinator 发生变化,Client 会标记旧 Coordinator 不可用并重新查找。

7.2 JoinGroup:汇报成员能力和订阅#

每个 Consumer 向 Coordinator 发送:

Member ID / Instance ID
Subscription
支持的 Assignment Strategies
Protocol Metadata
Rebalance Timeout

Coordinator 收集当前 Group 成员。

Classic Protocol 会选择一个成员作为:

Group Leader

注意:

Group Leader 是 Consumer Client,不是 Broker Leader。

7.3 Client-side Assignor 计算全组分配#

Coordinator 把成员和 Subscription 信息交给 Group Leader。

Group Leader 在客户端运行 Assignor:

GroupSubscription
Range / RoundRobin / Sticky ...
GroupAssignment

得到:

Member A → P0 P1
Member B → P2 P3

这就是 Classic Protocol 为什么把较多协调复杂性放在 Client。

7.4 SyncGroup:提交并分发 Assignment#

Group Leader 通过 SyncGroup 把每个成员的 Assignment 发给 Coordinator。

其他成员也发送 SyncGroup,等待 Coordinator 返回自己的 Assignment。

成功后 Group 进入:

Stable

Consumer 获得 Partition,初始化 Position,开始 Fetch。

7.5 Generation 的作用#

一次成功 Rebalance 会形成新的:

Generation

Offset Commit 等 Group 操作携带成员和 Generation 信息。

如果一个旧成员在 Rebalance 后继续提交:

旧 Generation
CommitFailed / Fenced

这防止已经失去 Ownership 的 Consumer 覆盖新 Owner 的进度。

7.6 为什么 Classic Rebalance 很贵#

一次成员变化可能导致:

所有成员进入 Join
等待成员集合稳定
Group Leader 重算 Assignment
所有成员 Sync
执行 Revoke / Assign Callback
重新初始化 Position 和 Fetch

它包含:

  • 全组同步屏障;
  • Partition Ownership 迁移;
  • 本地缓存丢弃;
  • Offset Commit;
  • 状态重建;
  • 外部资源初始化。

因此 Rebalance 的代价不能只用“几次网络请求”衡量。


8. Classic Group Coordinator 状态机#

Classic Group 状态机

经典 Group 状态通常包括:

Empty
PreparingRebalance
CompletingRebalance
Stable
Dead

8.1 Empty#

Group 没有活跃成员,但可能仍保存 Committed Offset。

首个成员加入后进入 Rebalance。

8.2 PreparingRebalance#

Coordinator 正在等待成员加入。

触发来源:

  • 新成员;
  • 成员离开;
  • Session Timeout;
  • Subscription 变化;
  • Topic Partition 数变化;
  • 显式 Rebalance;
  • 协议或 Assignor 变化。

8.3 CompletingRebalance#

成员集合已经确定,正在等待 Assignment 同步。

此时:

  • Group Leader 计算结果;
  • 成员等待 SyncGroup;
  • Assignment 尚未完全生效。

8.4 Stable#

成员拥有稳定 Assignment,正常:

Heartbeat
Fetch
Commit

8.5 Dead#

Group Coordinator 不再服务该 Group,例如内部 Partition 迁移或 Coordinator 关闭过程中。

8.6 状态机如何帮助排错#

如果 Commit 抛出:

RebalanceInProgressException

不是简单“网络不好”,而是 Group 正处于 Ownership 未稳定状态。

如果 Group 长时间:

PreparingRebalance

应该检查:

  • 是否有成员持续 Join / Leave;
  • 是否有实例无法在 Rebalance Timeout 内重新加入;
  • Callback 是否耗时;
  • Coordinator 是否不稳定;
  • Assignor 配置是否不兼容。

Kafka 4.3 Monitoring 文档为 Classic Group 暴露:

PreparingRebalance
CompletingRebalance
Empty
Stable
Dead

等状态计数指标。5


9. Partition Assignor:分配不是简单平均除法#

Kafka 必须决定:

Members × Subscriptions × Partitions
Assignment

不同 Assignor 优化目标不同。

Classic Assignor 对比

9.1 RangeAssignor#

按 Topic 分别分配连续 Partition 区间。

例如两个 Topic:

Topic A: A0 A1 A2
Topic B: B0 B1 B2
Consumers: C1 C2

可能:

C1 → A0 A1 B0 B1
C2 → A2 B2

优点:

  • 简单;
  • 同 Topic Partition 连续;
  • 适合某些局部性。

缺点:

  • 多 Topic 下可能产生倾斜;
  • 默认使用 Eager Rebalance。

Kafka 4.3 Classic Client 默认 Assignor 列表第一个仍是 Range,因此默认实际使用 Range。2

9.2 RoundRobinAssignor#

把所有已订阅 Partition 轮询分给成员。

优点:

  • 在 Subscription 一致时总体更均衡。

缺点:

  • Subscription 不一致时结果复杂;
  • Assignment 变化时可能移动大量 Partition;
  • 默认是 Eager。

9.3 StickyAssignor#

目标:

尽量平衡
+
尽量保留原 Assignment

减少 Rebalance 时 Partition 迁移。

但它仍是 Eager Rebalance Protocol:

成员先撤销全部 Ownership
再获得新 Assignment

所以“Sticky”主要描述分配结果稳定性,不等于增量 Ownership 迁移。

9.4 CooperativeStickyAssignor#

它使用 Sticky 分配思想,并启用 Cooperative Rebalance。

目标:

只撤销确实需要移动的 Partition。

Kafka 4.3 默认 Assignor 列表为:

RangeAssignor
CooperativeStickyAssignor

这样可以通过滚动升级移除 Range,把所有 Consumer 切换到 CooperativeSticky。2

9.5 Assignor 选择必须考虑什么#

不要只比较“均不均匀”。

至少评估:

Topic 数量
Subscription 是否一致
Partition 数量
Assignment 稳定性
Stateful Resource 初始化成本
Partition 迁移代价
Rebalance 停顿容忍度
滚动升级兼容

例如 Consumer 为每个 Partition 建立本地缓存或数据库连接时,Sticky 的价值会明显增加。


10. Eager 与 Cooperative:Ownership 如何迁移#

Eager 与 Cooperative Ownership 迁移

10.1 Eager Rebalance#

Eager 模型:

所有成员撤销全部 Partition
全组重新分配
成员获得新 Assignment

即使最终只有一个 Partition 需要从 A 移到 C,A 和 B 也可能都暂时停止所有 Partition。

这就是常说的:

Stop The World Rebalance

更准确地说,是:

Group 内 Partition Ownership 出现全量撤销窗口。

10.2 Cooperative Rebalance#

Cooperative 模型:

成员保留无需移动的 Partition
仅撤销存在冲突的 Partition
目标成员在旧 Owner 撤销后获取
通过多轮增量收敛

核心安全条件:

同一 Partition 不能同时被两个成员拥有。

因此 Cooperative 不是“一步直接覆盖 Assignment”,而是:

Revoke
→ 确认释放
→ Assign

10.3 Cooperative 为什么可能需要多轮#

假设:

C1 当前拥有 P0 P1
C2 当前拥有 P2 P3
新增 C3

目标:

C1 → P0
C2 → P2
C3 → P1 P3

第一轮先让:

C1 revoke P1
C2 revoke P3

第二轮 C3 才能安全获得。

因此 Cooperative 减少停顿范围,但可能通过多轮 Rebalance 收敛。

10.4 回调语义也不同#

Eager#

onPartitionsRevoked 通常收到旧 Assignment 全集。

Cooperative#

onPartitionsRevoked 通常只收到本轮需要转移的 Partition。

业务代码不能继续假设:

每次 Rebalance 都会 Revoke 全部 Partition。


11. Kafka 4.x 新 Consumer Rebalance Protocol#

Kafka 4.0 起,KIP-848 新 Consumer Protocol 已经 GA。1

Classic 与新 Consumer Protocol

启用:

group.protocol=consumer

当前默认仍是:

group.protocol=classic

11.1 为什么重新设计协议#

Classic Protocol 的关键限制包括:

  • Assignment 在 Client Group Leader 上计算;
  • JoinGroup / SyncGroup 形成全组同步屏障;
  • 大 Group 的 Subscription Metadata 很重;
  • 一个成员变化容易影响全组;
  • Client 配置控制 Heartbeat、Session 和 Assignor;
  • Broker 很难独立增量推进每个成员状态。

新协议目标:

更大 Consumer Group
更短 Rebalance
Broker 端集中 Assignment
Fully Incremental
简化 Client
改进线程模型

11.2 Broker-side Assignor#

新协议中,Assignment Strategy 由 Broker 控制:

group.consumer.assignors=uniform,range

Kafka 4.3 默认提供:

uniform
range

默认优先:

uniform

Consumer 可以通过:

group.remote.assignor=range

请求 Broker 端可用 Assignor。1

这与 Classic 的:

partition.assignment.strategy

完全不同。

11.3 Member Epoch#

Classic Protocol 常围绕:

Generation

管理全组版本。

新协议更强调:

Member Epoch
Target Assignment
Current Assignment

每个成员通过 Heartbeat 持续汇报自己的状态。

Broker 返回:

  • 当前成员 Epoch;
  • 目标 Assignment;
  • 需要 Revoke / Assign 的变化。

成员独立进行 Reconciliation。

11.4 Fully Incremental Reconciliation#

新协议不再要求:

所有成员同时停下
所有成员同时 Join
所有成员同时 Sync

而是:

Coordinator 更新目标 Assignment
受影响成员逐步 Reconcile
未受影响成员继续处理
Group 最终收敛

官方文档明确指出,新协议不再依赖 Global Synchronization Barrier。1

11.5 新 Consumer Group 状态#

Kafka 4.3 Monitoring 中,新 Consumer Group 主要状态为:

Empty
Assigning
Reconciling
Stable
Dead

其中:

Assigning#

Broker 正在计算或更新目标分配。

Reconciling#

成员当前 Assignment 正在向目标 Assignment 收敛。

这比 Classic 的 Preparing / Completing 更直接表达增量协调过程。5

11.6 Heartbeat 与 Session 配置移动到 Broker#

使用新协议时,以下 Client 配置不再支持:

heartbeat.interval.ms
session.timeout.ms
partition.assignment.strategy
enforceRebalance()

改由 Broker 配置:

group.consumer.heartbeat.interval.ms
group.consumer.session.timeout.ms
group.consumer.assignors

统一管理 Group 的成员检测和 Assignment 策略。1

11.7 新协议并不消灭所有 Rebalance 成本#

即使 Fully Incremental,以下工作仍可能发生:

  • Partition Revoke Callback;
  • 未处理 Record 清理;
  • Offset Commit;
  • 新 Partition Position 初始化;
  • State / Cache 初始化;
  • 外部连接迁移;
  • Hot Partition 转移;
  • 大量成员同时滚动发布。

新协议减少协调放大效应,但业务自身的 Ownership Migration 成本仍然存在。

11.8 Classic Cooperative 与新 Consumer Protocol 不要混淆#

CooperativeStickyAssignor
→ Classic Protocol 上的 Cooperative Ownership 迁移
group.protocol=consumer
→ 新的 Broker-driven Consumer Protocol

两者都强调 Incremental,但协议结构不同。

11.9 升级策略#

Kafka 4.3 支持:

  • Group 为空时离线切换;
  • 满足兼容条件时滚动在线升级;
  • 反向滚动降级。

但正式迁移前必须验证:

所有 Client 版本
自定义 Assignor / Metadata
监控指标
Callback 语义
Broker group configs
回滚路径

不能只改一个配置就直接全量生产发布。


12. Heartbeat、Session 与 max.poll.interval.ms#

这是 Consumer 最容易混淆的三个时间概念。

Heartbeat、Session 与 max.poll 时间线

12.1 Heartbeat Interval#

Classic Protocol:

heartbeat.interval.ms=3000

表示成员向 Coordinator 发送 Heartbeat 的期望间隔。

通常应小于 session.timeout.ms,官方建议一般不超过其三分之一。2

新 Consumer Protocol 中该 Client 配置无效,由 Broker:

group.consumer.heartbeat.interval.ms

控制。

12.2 Session Timeout#

Classic Protocol:

session.timeout.ms=45000

如果 Coordinator 在该时间内没有收到 Heartbeat:

Member 被认为失效
移出 Group
Partition 重新分配

新协议由 Broker:

group.consumer.session.timeout.ms

控制。

Session Timeout 关注:

Consumer 进程 / 网络是否仍能维持 Group Session。

12.3 max.poll.interval.ms#

默认:

max.poll.interval.ms=300000

它关注:

应用是否仍在持续调用 poll,推进业务消费循环。

即使 Heartbeat Thread 仍能发送 Heartbeat,如果应用长时间不 poll:

进程活着
Heartbeat 正常
业务处理卡死
Position 不再推进

这属于 Livelock。

Kafka 用 max.poll.interval.ms 防止一个“活着但不工作”的 Consumer 永久占有 Partition。2

12.4 为什么业务处理会触发 Rebalance#

代码:

while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
slowRemoteApi(record); // 总处理时间 8 分钟
}
}

如果:

max.poll.interval.ms = 5 分钟

Consumer 无法及时再次 poll。

Kafka 会认为它不再具备 Ownership 能力。

即使它最终处理完成,Commit 可能失败:

Partition 已经分给新 Member
旧 Member Commit
→ CommitFailedException

12.5 max.poll.records 与处理时间#

可以估算:

最大单批处理时间
max.poll.records
×
单条最坏处理时间

如果:

500 records
×
1 second
=
500 seconds

已经超过 5 分钟。

此时可选:

  • 降低 max.poll.records
  • 提高 max.poll.interval.ms
  • 批量处理;
  • 优化外部依赖;
  • pause Partition 后交给受控 Worker;
  • 将业务工作与 Poll Thread 解耦,但必须维护每 Partition 顺序和 Commit 水位。

12.6 Static Membership#

配置:

group.instance.id=fulfillment-consumer-01

Consumer 成为静态成员。

短暂重启时,Coordinator 可以保留其 Member Identity 和 Assignment,减少瞬时 Rebalance。2

但 Static Membership 不是无限租约。

达到 max.poll.interval.ms 后,静态成员会停止发送 Heartbeat,Partition 在 Session Timeout 后重新分配。2

它适合:

  • 实例有稳定唯一 ID;
  • 滚动重启;
  • 短暂网络抖动;
  • Stateful Consumer。

风险:

  • 重复 Instance ID 会 Fencing;
  • Session Timeout 太大导致真正故障接管慢;
  • 容器环境必须正确生成稳定且唯一的 ID。

12.7 不要用无限增大 Timeout 掩盖问题#

把:

session.timeout.ms
max.poll.interval.ms

改得非常大,只会:

  • 延迟真实故障接管;
  • 延迟 Partition 释放;
  • 扩大 Lag;
  • 隐藏业务线程卡死;
  • 让发布和扩缩容变慢。

正确做法是先测量:

Poll 间隔分布
Batch Processing P99
Heartbeat RTT
Rebalance Duration
外部依赖 Timeout

再设置故障检测窗口。


13. Rebalance 到底由什么触发#

13.1 新成员加入#

扩容:

C1 C2
新增 C3

需要把部分 Partition 转移给 C3。

13.2 成员正常退出#

close() 或 LeaveGroup:

C2 离开
C2 的 Partition 重新分配

正常关闭通常比直接 kill -9 更快释放 Ownership。

13.3 Session Timeout#

进程宕机、网络中断或 Heartbeat 停止:

Session Timeout
Coordinator 移除成员
Rebalance

13.4 Max Poll Interval#

应用处理卡住:

poll 不再推进
成员失去消费资格
Ownership 转移

13.5 Subscription 变化#

Consumer 调用:

subscribe(newTopics);

或者 Pattern Subscription 匹配到新 Topic。

13.6 Topic Partition 数变化#

Topic 扩容后,Group 需要把新增 Partition 纳入 Assignment。

13.7 Assignor / Protocol 变化#

滚动修改:

partition.assignment.strategy
group.protocol

可能触发协议协商和重新分配。

13.8 Coordinator 或内部 Topic 变化#

__consumer_offsets Leader 变化、Coordinator 迁移或 Broker 故障会导致 Client 重新查找 Coordinator,并可能影响 Group 进度。


14. ConsumerRebalanceListener:Ownership 迁移的业务钩子#

使用:

consumer.subscribe(topics, new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
}
@Override
public void onPartitionsLost(Collection<TopicPartition> partitions) {
}
});

14.1 onPartitionsRevoked#

在可控 Ownership 移交前调用。

常见工作:

  • 提交已处理 Offset;
  • Flush Partition 本地状态;
  • 停止对应 Worker;
  • 释放 Partition 专属资源;
  • 保存外部 Checkpoint。

Callback 必须:

快速
有界
可观测
幂等

如果 Callback 执行分钟级远程 I/O,会拖慢整个 Rebalance。

14.2 onPartitionsAssigned#

新 Assignment 生效后调用。

常见工作:

  • 读取外部 Offset;
  • seek() 到自定义位置;
  • 初始化 Partition State;
  • 恢复缓存;
  • 启动对应 Worker。

如果自定义 Offset 存在外部系统,应在这里恢复 Position。

14.3 onPartitionsLost#

表示 Partition Ownership 已经丢失,应用不应再假设自己能安全提交或完成最终处理。

典型场景:

  • Session 已经过期;
  • Coordinator 已经把 Partition 分给其他成员;
  • 无法执行正常 Revoke 流程。

此时不要继续:

把旧 Member 的 Offset 当作最终权威提交

14.4 Callback 在哪里执行#

Rebalance Callback 通常在 Consumer 应用线程的 poll 流程中执行。

因此 Callback 的耗时直接影响:

poll latency
Rebalance duration
Group convergence
业务处理停顿

不要在 Callback 内启动不受控并发后立刻返回,否则可能在旧 Ownership 工作仍未结束时让新 Owner 开始处理。


15. Offset Commit:什么时候才算“消费完成”#

Offset Commit 数据路径

15.1 Commit 的对象是恢复位置#

假设成功处理:

P0 Offset 0..99

应提交:

P0 → 100

含义:

重启后从 100 开始。

Commit 并不会删除 Kafka 中的 Record。

Retention 仍由 Topic 配置决定。

15.2 Offset Commit 如何持久化#

流程:

Consumer
↓ OffsetCommitRequest
Group Coordinator
↓ Append
__consumer_offsets
↓ Replication
CommitResponse

官方实现文档说明,Coordinator 把 Offset Commit 追加到 __consumer_offsets Compacted Topic,并在达到内部 Topic 的复制确认条件后返回成功;Coordinator 同时维护 Offset Cache 加速查询。4

15.3 enable.auto.commit=true#

Kafka 4.3 默认:

enable.auto.commit=true
auto.commit.interval.ms=5000

自动提交的是 Consumer 已经交付 / 推进的 Position。

它并不知道你的业务是否真正成功。

危险场景:

poll 返回 500 条
Position 已推进
Auto Commit 发生
业务只处理 100 条后进程崩溃

重启可能从提交后的更远位置继续,造成未处理 Record 被跳过。

因此包含重要业务副作用时,通常应:

enable.auto.commit=false

并由应用控制提交。

15.4 提交过早:At-Most-Once 风险#

Commit Offset
执行数据库更新
进程崩溃

重启从新 Offset 继续,旧 Record 不再重放。

业务更新可能丢失。

15.5 提交过晚:At-Least-Once 重复#

业务成功
Commit 前进程崩溃
重启重新消费

Record 会再次处理。

因此企业应用通常采用:

处理成功后提交
+
业务幂等

15.6 commitSync()#

同步等待 Commit 成功或失败:

consumer.commitSync();

优点:

  • 错误直接返回;
  • 顺序清晰;
  • 适合 Revoke、关闭或关键 Checkpoint。

缺点:

  • 阻塞 Poll Thread;
  • Coordinator 慢时增加处理延迟;
  • 每批都同步 Commit 会降低吞吐。

更安全的显式提交:

Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(tp, new OffsetAndMetadata(lastProcessedOffset + 1));
consumer.commitSync(offsets);

15.7 commitAsync()#

consumer.commitAsync(offsets, callback);

优点:

  • 不阻塞等待响应;
  • 更适合高频提交。

风险:

Commit 100 发出
Commit 200 发出
Commit 200 成功
Commit 100 延迟失败 / 回调

如果应用在 Callback 中盲目重试 100,可能把恢复点倒退。

需要维护:

Monotonic Commit Watermark

只允许提交进度向前。

15.8 常见组合#

业务循环中:

周期性 commitAsync

Rebalance / Close 时:

最后一次 commitSync

这样兼顾吞吐和退出确定性。

但仍要检查当前 Assignment 与处理水位。

15.9 多线程处理时提交什么#

假设 Poll Thread 把同一 Partition 的 Record 交给 Worker:

Offset 100 完成
Offset 101 仍处理中
Offset 102 完成

不能提交:

103

否则 101 失败后无法重放。

可提交的水位是:

从旧 Position 开始,已经连续完成的最大前缀的下一 Offset。

即:

Completed: 100, 102
Gap: 101
Safe Commit = 101

这就是多线程 Consumer 最难的部分之一。

15.10 Commit 不能解决外部副作用 Exactly Once#

Kafka Record
更新 MySQL
Commit Kafka Offset

MySQL Transaction 与 Kafka Offset Commit 是两个系统。

中间崩溃仍可能重复。

第五篇和第六篇会继续讨论:

  • Idempotent Consumer;
  • Unique Key;
  • Outbox;
  • Kafka Transaction;
  • Read-Process-Write EOS 边界。

16. Position 初始化、Offset Reset 与 Replay#

16.1 Assignment 后从哪里开始#

Consumer 获得 Partition 后:

存在 Committed Offset
→ 从 Committed Position 开始
不存在或 Offset 已过期
→ 使用 auto.offset.reset

16.2 auto.offset.reset#

Kafka 4.3 支持:

earliest
latest
by_duration:<ISO-8601 duration>
none

例如:

auto.offset.reset=earliest

从当前最早可用位置开始。

auto.offset.reset=latest

从当前最新位置开始,只消费后续 Record。

auto.offset.reset=by_duration:PT6H

尝试从当前时间向前六小时对应的位置开始。2

auto.offset.reset=none

没有有效 Committed Offset 时直接失败,适合不允许静默跳转的关键系统。

16.3 Offset 为什么会失效#

  • Retention 删除旧 Segment;
  • Group 长期不活跃导致 Offset 过期;
  • 管理员 Reset Offset;
  • Topic 重建;
  • Partition 发生变化;
  • 外部错误提交。

16.4 seek():移动当前 Position#

consumer.seek(tp, offset);

改变当前 Consumer 的读取 Position。

它不会自动修改 Committed Offset。

所以:

seek(100)
从 100 重读
如果未 Commit
重启仍从原 Committed Offset

16.5 Replay 的两种方式#

在线 Consumer 内 seek#

适合:

  • 小范围调试;
  • 自定义恢复;
  • Rebalance Listener 中外部 Offset 初始化。

管理工具 Reset Group Offset#

Terminal window
kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group order-fulfillment \
--topic order-events \
--reset-offsets \
--to-earliest \
--execute

适合整个 Group 的运维回放。

生产执行前必须:

  • 停止或隔离 Group;
  • 先使用 --dry-run
  • 明确业务幂等;
  • 评估回放流量;
  • 记录变更审计。

17. 慢消费、Pause 与多线程处理#

17.1 Lag 上升的基本模型#

设:

Producer Rate = λ
Consumer Processing Rate = μ

当:

λ > μ

Lag 会持续增长。

但根因可能在不同层:

Partition 不足
Consumer 不足
单 Partition 热点
Fetch 慢
反序列化慢
业务处理慢
DB / API 慢
Commit 慢
频繁 Rebalance

17.2 pause() 不会释放 Ownership#

consumer.pause(partitions);

表示暂时不从这些 Partition 向应用交付新 Record。

它不会:

  • 退出 Group;
  • 触发 Rebalance;
  • 把 Partition 给其他 Consumer。

适用于:

  • Worker Queue 已满;
  • 下游限流;
  • 单 Partition 暂停;
  • 保持 Poll / Heartbeat,同时停止获取更多业务工作。

恢复:

consumer.resume(partitions);

17.3 Poll Thread + Worker Pool#

常见架构:

Poll Thread
按 Partition 分发
Partition Worker / Ordered Queue
业务处理
Completed Offset Tracker
Commit Safe Watermark

关键规则:

  1. 同一 Partition 的业务顺序必须受控;
  2. Poll Thread 仍需持续 poll;
  3. Worker Queue 满时 pause;
  4. 只提交连续完成前缀;
  5. Revoke 时停止接收新任务并等待 / 取消旧任务;
  6. Ownership Lost 后不再提交旧结果。

这不是简单:

records.parallelStream().forEach(this::handle);

17.4 什么时候应该扩 Consumer#

同时满足:

存在未利用 Partition 并行度
处理瓶颈可通过更多实例分摊
下游系统也能承受扩容流量
Assignment 和 Rebalance 成本可控

如果:

Partitions = Consumers

继续加 Consumer 没有直接效果。

如果只有 P7 是 Hot Partition,加十个 Consumer 也无法拆分 P7 内部顺序流。


18. Consumer Metrics:从指标判断卡在哪一层#

18.1 Fetch 层#

关注:

fetch-rate
fetch-latency-avg / max
fetch-size-avg / max
records-per-request-avg
bytes-consumed-rate
records-consumed-rate

判断:

  • Fetch 是否频繁但 Batch 很小;
  • Broker 响应是否变慢;
  • 网络是否成为瓶颈;
  • Consumer 是否真正有数据流入。

18.2 Lag 层#

records-lag
records-lag-max
records-lead

必须按 Partition 观察。

总 Lag 可能掩盖:

一个 Hot Partition
+
大量空闲 Partition

18.3 Poll 与处理层#

last-poll-seconds-ago
time-between-poll-avg / max
poll-idle-ratio-avg
records-consumed-rate

如果:

time-between-poll-max
接近 max.poll.interval.ms

说明业务处理已经接近失去 Ownership。

18.4 Coordinator 与 Rebalance#

Classic / Consumer 指标关注:

assigned-partitions
heartbeat-rate
heartbeat-response-time-max
last-heartbeat-seconds-ago
rebalance-total
rebalance-rate-per-hour
last-rebalance-seconds-ago
failed-rebalance-total

Broker 侧 Kafka 4.3 还暴露:

consumer-group-count{state}
consumer-group-rebalance-rate
classic group completed rebalance rate
offset-commit-rate

等 Group Coordinator 指标。5

18.5 Commit#

commit-rate
commit-latency-avg
commit-latency-max
commit-failed-rate

Commit 慢可能意味着:

  • Coordinator Broker 压力;
  • __consumer_offsets 副本问题;
  • 网络抖动;
  • Group 正在 Rebalance;
  • Commit 过于频繁。

18.6 指标必须组合解释#

例如:

Lag ↑
records-consumed-rate ↓
fetch-latency 正常
time-between-poll ↑

更可能是业务处理慢,而不是 Broker Fetch。

Lag ↑
fetch-latency ↑
heartbeat 正常
业务 CPU 低

更可能是 Broker、网络或 Quota。


19. Consumer 故障诊断树#

Consumer 故障诊断树

19.1 症状:Lag 持续上升#

按顺序检查:

Ownership#

Group 是否 Stable?
Rebalance Rate 是否异常?
Assigned Partitions 是否减少?
是否存在空闲 Consumer?

Fetch#

Leader 是否可用?
Fetch Latency 是否升高?
是否有大 Batch 卡住?
Broker Quota / TLS 是否异常?

Processing#

DB / HTTP P99
Worker Queue
GC
CPU
Hot Partition

Commit#

Committed Lag 高但 Position Lag 低?
Commit 是否过慢或失败?

19.2 症状:重复消费突然增多#

检查:

  • 业务成功后 Commit 前崩溃;
  • Commit 失败;
  • Rebalance 期间未提交;
  • Async Commit Callback 重试倒退;
  • Offset 被人工 Reset;
  • Consumer 被频繁 Kill;
  • 外部系统不幂等;
  • Worker 完成顺序与 Commit 水位错误。

19.3 症状:频繁 Rebalance#

检查:

实例是否频繁重启?
max.poll.interval 是否过小?
Callback 是否太慢?
session timeout 是否过小?
网络是否丢 Heartbeat?
容器 group.instance.id 是否冲突?
Assignor 列表是否一致?
Topic Pattern 是否频繁匹配变化?

19.4 症状:Consumer 存活但不再处理#

可能是:

  • Poll Thread 死锁;
  • Worker Queue 阻塞;
  • Deserializer 卡住;
  • Rebalance Callback 卡住;
  • 所有 Partition 被 Pause;
  • wakeup() 状态处理错误;
  • read_committed 被 Open Transaction 阻挡;
  • Consumer 没有 Assignment。

19.5 症状:增加 Consumer 后吞吐不升#

检查:

Partition 数是否足够?
是否只有单个 Hot Partition?
Assignor 是否导致倾斜?
下游 DB 是否已经饱和?
Broker / Network 是否瓶颈?
新 Consumer 是否真正获得 Assignment?

20. 源码阅读地图#

Kafka 4.x Consumer 同时存在 Classic 与新 Async Runtime。

阅读时必须先确认:

group.protocol=classic

还是:

group.protocol=consumer

20.1 公共入口#

阅读:

KafkaConsumer
Consumer
ConsumerDelegate
ConsumerConfig
SubscriptionState
ConsumerMetadata

重点问题:

  1. KafkaConsumer 如何选择 Delegate?
  2. subscribe、assign、poll、commit、seek 如何委托?
  3. SubscriptionState 保存哪些 Position 与 Assignment?
  4. Thread Ownership 如何检查?
  5. wakeup 如何跨线程生效?

20.2 Classic Client 路线#

建议阅读:

ClassicKafkaConsumer
ConsumerCoordinator
AbstractCoordinator
Heartbeat
HeartbeatThread
ConsumerNetworkClient
Fetcher / FetchCollector
SubscriptionState

关注:

poll()#

updateAssignmentMetadataIfNeeded
pollForFetches
processBackgroundEvents / callbacks
return ConsumerRecords

具体方法名会随版本演进,应固定 Kafka 4.3 Tag。

ConsumerCoordinator#

ensureCoordinatorReady
ensureActiveGroup
onJoinPrepare
performAssignment
onJoinComplete
commitOffsetsSync / Async
maybeAutoCommit

AbstractCoordinator#

FindCoordinator
JoinGroup
SyncGroup
Heartbeat Thread
Generation / Member ID
Rejoin Need

20.3 新 Async Consumer 路线#

阅读:

AsyncKafkaConsumer
ApplicationEventHandler
ConsumerNetworkThread
NetworkClientDelegate
ConsumerMembershipManager
HeartbeatRequestManager
CommitRequestManager
FetchRequestManager
OffsetsRequestManager
BackgroundEventProcessor

关注:

  1. Application API 如何转换为 Event?
  2. ConsumerNetworkThread 如何 poll 网络?
  3. Heartbeat Response 如何推进 Member Epoch?
  4. Target Assignment 与 Current Assignment 在哪里 Reconcile?
  5. Rebalance Callback 如何回到应用线程?
  6. Commit Future 如何跨线程完成?
  7. Fetch Buffer 如何交付给 poll?

20.4 Broker Group Coordinator 路线#

阅读:

GroupCoordinatorService
GroupCoordinatorShard
GroupMetadataManager
OffsetMetadataManager
ClassicGroup
ConsumerGroup
ConsumerGroupMember
TargetAssignmentBuilder

关注:

Classic Group#

PreparingRebalance
CompletingRebalance
Stable
Generation
Protocol
Leader

Consumer Group#

Group Epoch
Member Epoch
Target Assignment
Current Assignment
Assigning / Reconciling / Stable

Offset#

OffsetCommit
__consumer_offsets Record
Offset Cache
Offset Fetch
Expiration / Tombstone

20.5 推荐阅读顺序#

KafkaConsumer Javadoc
SubscriptionState / ConsumerRecords
ClassicKafkaConsumer.poll
FetchCollector / Fetcher
ConsumerCoordinator
AbstractCoordinator
Classic Group Broker State
然后:
Consumer Rebalance Protocol 官方文档
AsyncKafkaConsumer
Application Event / ConsumerNetworkThread
Membership / Heartbeat Managers
ConsumerGroup Broker State

不要一开始直接读 Group Coordinator 数万行代码。

先建立:

Position
Membership
Assignment
Commit

四个状态域。

20.6 源码阅读记录模板#

每条链固定记录:

入口 API
运行线程
输入状态
发出的 Request / Event
Broker 状态变化
Client 本地状态变化
回调 / 返回值
错误和重试边界

例如 Rebalance:

入口:poll()
线程:Application + Heartbeat / Network
输入:Subscription / Member Epoch
请求:JoinGroup 或 ConsumerGroupHeartbeat
Broker:Target Assignment
Client:Revoke / Assign / Position Init
回调:ConsumerRebalanceListener
错误:Coordinator Change / Fencing / Timeout

21. 可运行实验:亲手观察 Consumer Group 与 Rebalance#

下面使用单节点 Kafka 4.3.0 做本地实验。

环境定位:学习和协议观察,不代表生产部署建议。

21.1 目录结构#

kafka-consumer-lab/
├── docker-compose.yml
├── pom.xml
└── src/main/java/lab/ConsumerLab.java

21.2 Docker Compose#

services:
kafka:
image: apache/kafka:4.3.0
container_name: kafka-ch04
hostname: kafka
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
# 允许同时实验 Classic 与 Consumer Protocol
KAFKA_GROUP_CONSUMER_HEARTBEAT_INTERVAL_MS: 5000
KAFKA_GROUP_CONSUMER_SESSION_TIMEOUT_MS: 45000
volumes:
- kafka_ch04_data:/var/lib/kafka/data
volumes:
kafka_ch04_data:

启动:

Terminal window
docker compose up -d
docker logs -f kafka-ch04

21.3 创建六 Partition Topic#

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic order-events \
--partitions 6 \
--replication-factor 1

验证:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic order-events

21.4 发送带 Key 的测试消息#

Terminal window
python3 - <<'PY' | docker exec -i kafka-ch04 \
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic order-events \
--property parse.key=true \
--property key.separator=:
for i in range(120):
order_id = f"order-{i % 12:02d}"
print(f"{order_id}:event-{i:03d}")
PY

同一个 orderId 会稳定映射到同一 Partition,方便观察 Partition 内顺序与 Group Assignment。

21.5 Maven 配置#

<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>lab</groupId>
<artifactId>kafka-consumer-lab</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.release>17</maven.compiler.release>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>4.3.0</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>2.0.17</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.codehaus.mojo</groupId>
<artifactId>exec-maven-plugin</artifactId>
<version>3.5.0</version>
<configuration>
<mainClass>lab.ConsumerLab</mainClass>
</configuration>
</plugin>
</plugins>
</build>
</project>

21.6 ConsumerLab.java#

package lab;
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.WakeupException;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.*;
import java.util.concurrent.atomic.AtomicBoolean;
public final class ConsumerLab {
private static final AtomicBoolean CLOSED = new AtomicBoolean(false);
private ConsumerLab() {
}
public static void main(String[] args) {
String instance = arg(args, "--instance", "consumer-1");
String group = arg(args, "--group", "order-lab");
String protocol = arg(args, "--protocol", "classic");
long sleepMs = Long.parseLong(arg(args, "--sleep-ms", "100"));
boolean manualCommit =
Boolean.parseBoolean(arg(args, "--manual-commit", "true"));
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, group);
props.put(ConsumerConfig.CLIENT_ID_CONFIG, instance);
props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, protocol);
props.put(
ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
props.put(
ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName()
);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(
ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,
Boolean.toString(!manualCommit)
);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "20");
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "30000");
// 只有 Classic Protocol 支持这些客户端配置。
if ("classic".equalsIgnoreCase(protocol)) {
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "3000");
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000");
props.put(
ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
CooperativeStickyAssignor.class.getName()
);
}
try (KafkaConsumer<String, String> consumer =
new KafkaConsumer<>(props)) {
Runtime.getRuntime().addShutdownHook(
new Thread(() -> {
CLOSED.set(true);
consumer.wakeup();
}, "consumer-shutdown")
);
consumer.subscribe(
List.of("order-events"),
new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(
Collection<TopicPartition> partitions
) {
System.out.printf(
"[%s] revoked %s%n",
instance,
partitions
);
}
@Override
public void onPartitionsAssigned(
Collection<TopicPartition> partitions
) {
System.out.printf(
"[%s] assigned %s%n",
instance,
partitions
);
}
@Override
public void onPartitionsLost(
Collection<TopicPartition> partitions
) {
System.out.printf(
"[%s] lost %s%n",
instance,
partitions
);
}
}
);
while (!CLOSED.get()) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(1000));
Map<TopicPartition, OffsetAndMetadata> processed =
new HashMap<>();
for (ConsumerRecord<String, String> record : records) {
System.out.printf(
"[%s] %s-%d offset=%d key=%s value=%s%n",
instance,
record.topic(),
record.partition(),
record.offset(),
record.key(),
record.value()
);
Thread.sleep(sleepMs);
processed.put(
new TopicPartition(
record.topic(),
record.partition()
),
new OffsetAndMetadata(record.offset() + 1)
);
}
if (manualCommit && !processed.isEmpty()) {
consumer.commitSync(processed);
}
}
} catch (WakeupException e) {
if (!CLOSED.get()) {
throw e;
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
private static String arg(
String[] args,
String name,
String defaultValue
) {
for (int i = 0; i < args.length - 1; i++) {
if (name.equals(args[i])) {
return args[i + 1];
}
}
return defaultValue;
}
}

21.7 实验一:同 Group 三个 Consumer#

打开三个终端。

终端一:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance c1 --group group-a --protocol classic"

终端二:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance c2 --group group-a --protocol classic"

终端三:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance c3 --group group-a --protocol classic"

观察:

六个 Partition 如何分给三个 Consumer
新增 Consumer 时哪些 Partition 被 Revoke
CooperativeSticky 是否保留已有 Assignment
每个 Partition 是否只出现在一个 Member 的 Assignment 中

停止其中一个 Consumer,再观察其 Partition 如何被其他成员接管。

21.8 实验二:不同 Group 独立消费#

另开终端:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance risk-1 --group risk-group --protocol classic"

它会拥有另一套独立 Offset,并从自己的恢复位置读取消息。

验证:

同 Group
→ Partition 分工
不同 Group
→ 各自拥有完整消费视图

21.9 实验三:观察 Group、Member 与 Lag#

查看 Group Offset:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group group-a

重点查看:

TOPIC
PARTITION
CURRENT-OFFSET
LOG-END-OFFSET
LAG
CONSUMER-ID
HOST
CLIENT-ID

查看成员和 Assignment:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group group-a \
--members \
--verbose

查看状态:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group group-a \
--state

21.10 实验四:制造慢消费与 max.poll.interval.ms 超时#

启动慢 Consumer:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance slow-1 --group slow-group --protocol classic --sleep-ms 2000"

当前参数:

max.poll.records = 20
单条 sleep = 2 秒
单批最坏处理时间约 40 秒
max.poll.interval.ms = 30 秒

预期现象:

  • Consumer 不能在 Poll Interval 内再次调用 poll;
  • Partition Ownership 被收回;
  • Commit 可能抛出 CommitFailedException
  • Group 发生新的 Rebalance。

分别尝试:

降低 max.poll.records
减少处理耗时
提高 max.poll.interval.ms

不要只比较“是否还报错”,还要比较:

故障接管时间
单批吞吐
Lag 增长速度
Commit 频率

21.11 实验五:切换到新 Consumer Protocol#

启动第一个成员:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance n1 --group new-group --protocol consumer"

再启动:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance n2 --group new-group --protocol consumer"
Terminal window
mvn -q exec:java \
-Dexec.args="--instance n3 --group new-group --protocol consumer"

观察:

  • 新 Group 能否正常获得 Assignment;
  • 客户端不再配置 Classic Heartbeat / Session 参数;
  • 滚动加入和退出成员时 Assignment 如何变化;
  • Group State 是否经过 Assigning、Reconciling、Stable;
  • 未发生迁移的 Partition 是否仍能继续处理。

查看状态:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe \
--group new-group \
--state

21.12 实验六:Offset Reset 与 Replay#

先停止 group-a 全部 Consumer。

Dry Run:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group group-a \
--topic order-events \
--reset-offsets \
--to-earliest \
--dry-run

确认后执行:

Terminal window
docker exec kafka-ch04 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--group group-a \
--topic order-events \
--reset-offsets \
--to-earliest \
--execute

重新启动 Consumer,观察消息 Replay。

正式生产环境执行 Reset 前必须:

停止或隔离正在运行的 Group
确认业务幂等
评估回放流量
先 Dry Run
保留变更审计
准备回滚位置

21.13 实验七:比较自动提交与手动提交#

自动提交:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance auto-1 --group auto-group --protocol classic --manual-commit false --sleep-ms 1000"

在处理过程中强制终止进程,再重启。

手动提交:

Terminal window
mvn -q exec:java \
-Dexec.args="--instance manual-1 --group manual-group --protocol classic --manual-commit true --sleep-ms 1000"

比较:

哪些 Offset 被重复
哪些 Record 可能被跳过
CURRENT-OFFSET 与实际业务完成位置是否一致

21.14 实验回验清单#

你应能够观察到:

  • 同 Group 内 Partition 不重复分配;
  • 不同 Group 独立消费;
  • 新成员加入导致 Assignment 变化;
  • Consumer 异常退出后 Partition 被接管;
  • 慢处理触发 Max Poll Timeout;
  • Commit 的 Offset 是下一读取位置;
  • Classic 和 Consumer Protocol 都能运行;
  • Group 工具可以显示 State、Member、Offset 与 Lag;
  • Offset Reset 可以触发 Replay;
  • 自动提交与业务完成状态可能存在偏差。

当前文档中的 Docker、Maven 和 Java 代码已完成静态检查。不同 Docker Desktop、JDK、Maven 和 Kafka 4.3.x 补丁版本可能需要调整环境变量或命令。本生成环境没有实际启动 Docker 与 Maven 进行端到端回验。


22. 常见错误认知#

22.1 “poll() 一次就是向 Broker 请求一次”#

错误。

Consumer 具有:

并行 Fetch
Completed Fetch
本地缓存
分批 poll 返回

一次 poll 可能完全消费本地缓存,也可能只推进网络和协调状态而不返回 Record。

22.2 “Heartbeat 正常就不会 Rebalance”#

错误。

Heartbeat 证明:

成员进程和网络仍能维持 Session

但应用长时间不 poll 仍会触发:

max.poll.interval.ms

Kafka 需要同时识别 Crash 与 Livelock。

22.3 “Committed Offset 是最后处理成功的 Offset”#

不准确。

它应该表示:

下一条需要读取的 Offset

成功处理 99 后提交 100。

22.4 “StickyAssignor 就是 Cooperative Rebalance”#

错误。

Sticky
→ Assignment 结果尽量稳定
Cooperative
→ Partition Ownership 增量迁移

StickyAssignor 仍使用 Eager Protocol;CooperativeStickyAssignor 才结合 Cooperative Rebalance。

22.5 “新 Consumer Protocol 默认已经启用”#

错误。

Kafka 4.3 客户端默认仍是:

group.protocol=classic

必须显式设置:

group.protocol=consumer

22.6 “Consumer 越多,吞吐一定越高”#

错误。

传统 Consumer Group 的有效并行度受:

Partition 数
Hot Partition
Broker 能力
下游能力

共同限制。

22.7 “把 Timeout 调大就能解决 Rebalance”#

通常只是延迟真实故障发现。

Timeout 调整必须建立在:

Poll Interval Distribution
Processing P99
Heartbeat RTT
Rebalance Duration

这些测量结果上。

22.8 “手动 Commit 就不会重复消费”#

错误。

业务成功后、Commit 成功前发生故障,仍会重新消费。

Manual Commit
Exactly Once Business Effect

22.9 “pause() 会把 Partition 交给其他消费者”#

错误。

pause() 只停止当前 Consumer 对 Partition 的数据交付,不释放 Ownership,也不会主动触发 Rebalance。

22.10 “Lag 高就应该立刻加 Consumer”#

错误。

必须先判断:

所有 Partition 都有 Lag?
只有一个 Hot Partition?
Group 是否 Stable?
Fetch 正常吗?
业务处理还是 Commit 慢?
Partition 数是否还有并行空间?

22.11 “Classic Cooperative 与 KIP-848 是同一种协议”#

错误。

CooperativeStickyAssignor
→ Classic Protocol 上的增量 Ownership 迁移
group.protocol=consumer
→ Broker-driven 新 Consumer Protocol

它们设计目标相近,但状态模型和请求协议不同。

22.12 “Kafka 保存 Offset,因此知道业务是否成功”#

错误。

Kafka 只知道应用提交了什么恢复位置。

它不知道:

MySQL 是否更新成功
HTTP 是否真正执行
扣款是否发生
邮件是否发送

业务副作用仍需幂等与事务边界设计。


23. 场景题#

23.1 场景一:八个 Partition,十二个 Consumer#

结论:

最多八个 Consumer 获得 Partition
至少四个 Consumer 没有有效 Assignment

进一步分析:

  • 是否应减少实例;
  • 是否为了容灾保留少量 Warm Standby;
  • Partition 是否真的需要扩容;
  • 扩 Partition 会不会破坏 Key 顺序假设。

23.2 场景二:Group 每分钟 Rebalance 一次#

排查顺序:

Group 当前 State
Member 启停记录
max.poll.interval.ms
Heartbeat / Session
Rebalance Callback Duration
group.instance.id
Assignor 配置一致性
Topic Pattern / Partition 变化
Coordinator 稳定性

不能直接把 session.timeout.ms 调大结束排查。

23.3 场景三:Lag 高但 Consumer CPU 很低#

可能原因:

  • DB 或 HTTP 阻塞;
  • Fetch Latency 高;
  • Consumer 没有 Assignment;
  • 所有 Partition 被 Pause;
  • 单个 Hot Partition;
  • Rebalance Storm;
  • Broker Quota;
  • read_committed 被 Open Transaction 的 LSO 阻挡;
  • Worker Queue 已满但 Poll Thread 仍空转。

23.4 场景四:业务偶尔跳过消息#

检查:

enable.auto.commit 是否开启?
是否处理前 Commit?
异步 Worker 未完成是否已提交?
是否直接提交 poll 返回的最大 Offset?
auto.offset.reset 是否为 latest?
是否人工 Reset 过 Group Offset?
是否错误处理反序列化异常?

23.5 场景五:业务出现重复扣款#

重复消费本身符合 At-Least-Once。

需要:

  • 业务幂等 Key;
  • 数据库唯一约束;
  • 状态机条件更新;
  • 处理结果表;
  • 正确 Commit 水位;
  • 外部支付接口幂等参数。

不能只调整 Kafka Consumer 参数。

23.6 场景六:滚动发布造成大面积停顿#

检查:

  • 是否使用 Eager Assignor;
  • Pod 是否同时终止;
  • 是否缺少 Graceful Close;
  • Session Timeout 是否导致接管过慢;
  • Callback 初始化是否太慢;
  • Consumer 是否维护昂贵本地 State;
  • 是否适合 CooperativeSticky;
  • 是否评估新 Consumer Protocol;
  • Static Membership 的 Instance ID 是否稳定且唯一。

23.7 场景七:增加 Consumer 后吞吐完全不变#

假设:

Topic = 6 Partitions
原 Consumer = 6
新增后 Consumer = 12

新增成员没有 Partition 可以获得。

另一个可能:

P0 占总流量 80%

即使其他 Partition 被均匀分配,P0 仍只有一个 Owner。

这时要解决的是:

Key Distribution / Hot Partition

而不是继续扩 Consumer。

23.8 场景八:Commit 偶尔失败,但业务已经成功#

可能流程:

业务处理成功
Group 在此刻 Rebalance
旧 Member 失去 Ownership
CommitFailedException

Record 会被新 Owner 再次读取。

因此业务必须幂等。

同时检查:

  • 单批处理时间;
  • Callback;
  • Group 稳定性;
  • Commit 是否只提交当前 Assignment;
  • Async Commit 是否乱序。

23.9 场景九:Consumer 没有报错,但一直拿不到数据#

检查:

是否真正 Assigned Partition?
Position 是否已经到可见上界?
是否所有 Partition 被 Pause?
Group ID 是否使用旧 Committed Offset?
auto.offset.reset 是否为 latest?
ACL 是否允许 Read?
read_committed 是否被事务阻挡?
Fetch 是否在重试 Leader / Metadata?

23.10 场景十:不同消费者处理同一订单时出现乱序#

先确定:

这些事件是否拥有相同 Key?
是否进入同一 Partition?
是否属于同一 Consumer Group?
业务是否把同一 Partition Record 并发分发到多个 Worker?
是否跨 Topic 合并顺序?
是否在 Retry / DLT 后改变顺序?

Kafka 只提供 Partition 内的日志顺序。

应用异步处理仍可能破坏最终业务顺序。


24. 本篇验收清单#

24.1 画图验收#

能够画出数据链:

Partition Leader
FetchRequest / FetchResponse
Completed Fetch
poll()
Deserialize
ConsumerRecords
Business Processing
Commit Safe Offset

能够画出 Classic 协调链:

FindCoordinator
→ JoinGroup
→ Client Leader Assign
→ SyncGroup
→ Stable

能够画出新 Consumer Protocol:

ConsumerGroupHeartbeat(Member Epoch)
→ Broker Target Assignment
→ Incremental Reconcile
→ Stable

24.2 Offset 验收#

能够明确区分:

Record Offset
Consumer Position
Committed Offset
Log End Offset
High Watermark
Last Stable Offset

并说明:

  • 哪些是 Consumer 本地状态;
  • 哪些是持久恢复状态;
  • 哪些是 Broker Log 边界;
  • read_committed 使用哪个可见上界。

24.3 Group 验收#

能够解释:

Group ID
Member ID
Static Instance ID
Generation
Member Epoch
Subscription
Assignment
Partition Ownership
Group Coordinator

以及这些概念分别属于 Classic 还是新 Consumer Protocol。

24.4 API 验收#

能够解释:

subscribe()
assign()
poll()
pause()
resume()
seek()
position()
committed()
commitSync()
commitAsync()
wakeup()
close()

每个 API 改变了哪一类状态。

24.5 配置验收#

可以把以下配置放回对应机制,而不是背参数表:

group.id
group.protocol
group.instance.id
partition.assignment.strategy
group.remote.assignor
heartbeat.interval.ms
session.timeout.ms
max.poll.interval.ms
max.poll.records
enable.auto.commit
auto.commit.interval.ms
auto.offset.reset
fetch.min.bytes
fetch.max.wait.ms
fetch.max.bytes
max.partition.fetch.bytes
isolation.level

24.6 Rebalance 验收#

面对 Rebalance,首先回答:

触发原因是什么?
Classic 还是 Consumer Protocol?
Group 当前状态是什么?
哪些 Partition 真正迁移?
Callback 是否耗时?
Offset 是否安全提交?
旧 Owner 是否还有未完成任务?
新 Owner 从哪里恢复?

24.7 Lag 验收#

面对 Lag,首先回答:

Lag 使用的左右坐标是什么?
所有 Partition 都高还是单个热点?
Group 是否 Stable?
Fetch 是否正常?
poll 是否及时?
业务处理是否阻塞?
Commit 是否推进?
Partition 数是否仍有并行空间?

24.8 多线程验收#

对于 Poll Thread + Worker Pool,可以说明:

  • 为什么同 Partition 不能无序并发;
  • 如何 pause / resume;
  • 如何计算连续完成前缀;
  • Revoke 时如何停止 Worker;
  • Lost 后为什么不能继续提交;
  • 如何控制 Worker Queue Backpressure。

24.9 源码验收#

至少能够说出两条 Client 路线。

Classic#

KafkaConsumer
→ ClassicKafkaConsumer
→ ConsumerCoordinator / AbstractCoordinator
→ Fetcher / FetchCollector
→ ConsumerNetworkClient

Consumer Protocol#

KafkaConsumer
→ AsyncKafkaConsumer
→ ApplicationEventHandler
→ ConsumerNetworkThread
→ Membership / Heartbeat / Commit / Fetch Managers

Broker 侧能够说出:

GroupCoordinatorService
→ GroupCoordinatorShard
→ GroupMetadataManager / OffsetMetadataManager
→ ClassicGroup 或 ConsumerGroup

24.10 实验验收#

实际操作并解释:

  1. 同 Group 启动三个 Consumer;
  2. 停止一个成员并观察接管;
  3. 不同 Group 独立消费;
  4. 使用命令查看 Member、Assignment 与 Lag;
  5. 制造 max.poll.interval.ms 超时;
  6. 比较 Classic 与 Consumer Protocol;
  7. Reset Offset 并 Replay;
  8. 比较自动与手动 Commit。

24.11 最终口述验收#

不看文章,用 45 分钟回答:

Consumer 调用 subscribe() 后,如何找到 Group Coordinator、加入 Group并获得 Partition Assignment;poll() 又如何推进 Fetch、本地缓存和 Position;业务处理完成后 Offset 如何提交到 __consumer_offsets;成员变化时 Classic Protocol 为什么需要 JoinGroup / SyncGroup,新 Consumer Protocol 又如何通过 Broker-side Assignor、Member Epoch 和 Incremental Reconciliation 降低全组协调成本。

如果只能回答:

消费者组会自动负载均衡,
消费者宕机后 Kafka 会 Rebalance。

说明仍然停留在使用层。


25. 全文收束:Consumer 是一个分布式 Ownership Runtime#

Kafka Consumer 的核心执行链可以收束为四层。

25.1 数据层#

Partition Leader
FetchRequest / Response
Completed Fetch
poll()
Deserialize
Business Processing

25.2 Position 层#

Record Offset
Consumer Position
Processed Watermark
Committed Offset
Failure Recovery

25.3 Membership 层#

Group ID
Group Coordinator
Member / Heartbeat
Session Liveness

25.4 Ownership 层#

Subscription
Assignment
Partition Ownership
Revoke / Assign
Incremental Reconciliation

最终可以得到一句更准确的定义:

Kafka Consumer Group 不是一组简单的消息读取线程,而是一个把 Partition Ownership、成员存活、读取位置和失败恢复统一起来的分布式协调 Runtime。

Classic Protocol 使用:

JoinGroup
+
SyncGroup
+
Client-side Assignor
+
Generation

维护全组 Assignment。

Kafka 4.x 新 Consumer Protocol 使用:

Broker-side Assignor
+
Member Epoch
+
Target Assignment
+
Incremental Reconciliation

减少全局屏障和协调放大。

但无论协议如何演进,应用仍要负责最关键的一段:

Record 已交付
业务真正成功
什么时候安全 Commit

Kafka 可以管理 Partition Ownership 和 Offset Recovery。

它不能自动知道:

数据库更新是否成功
HTTP 调用是否产生副作用
扣款是否真正完成
外部系统是否幂等

下一篇,我们会进入 Kafka 最核心的 Failure Model:

Producer ACK
Leader / Follower
Replication
ISR
High Watermark
Leader Failure
Duplicate
Idempotent Producer
Transaction
Exactly Once Boundary

回答:

分布式系统必然失败时,Kafka 如何定义“消息已经成功”,又如何在丢失、重复、乱序和跨 Partition 原子性之间建立可推导的可靠性模型?


参考资料#

  1. Apache Kafka 4.3 — Consumer Rebalance Protocol
  2. Apache Kafka 4.3 — Consumer and Share Consumer Configs
  3. Apache Kafka 4.3 — KafkaConsumer Javadoc
  4. Apache Kafka 4.3 — Distribution / Consumer Offset Tracking
  5. Apache Kafka 4.3 — Monitoring
  6. Apache Kafka 4.3 — Design
  7. Apache Kafka 4.3 — CooperativeStickyAssignor Javadoc
  8. Apache Kafka 4.3 — ConsumerPartitionAssignor Javadoc
  9. Apache Kafka 4.3 — Basic Kafka Operations
  10. Apache Kafka 4.3.0 Release Announcement
  11. Apache Kafka Source Repository
  12. Apache Kafka 4.3 — Upgrade Guide / Share Groups
Kafka Consumer Group:从 poll()、Offset 到 Rebalance 的分布式协调系统
https://jupiter-ws.cn/posts/backend/kafka/04_kafka_consumer_group_rebalance/
作者
Jupiter
发布于
2026-07-17
许可协议
CC BY-NC-SA 4.0