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 AssignmentRebalance CallbackOffset Fetch / CommitMetadata 与网络 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.ms、session.timeout.ms与max.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 AssignmentGroup 正在 Rebalance目标 Position 已经追到可见上界read_committed 被 LSO 阻挡Partition 被 pauseMetadata 或 Leader 暂时不可用
更重要的是,即便这一次没有返回 Record,poll() 仍可能在推进:
网络事件Group JoinAssignment 回调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 LivenessOwnership 链
Subscription ↓Assignment Calculation ↓Partition Revocation / Assignment ↓ReconciliationOffset 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 的 PositionPause 状态Pending CommitRebalance Callback本地 Fetch BufferGroup 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 Protocol
group.protocol=classic概念主线:
Application Thread ↓ClassicKafkaConsumer ↓ConsumerCoordinator / SubscriptionState / Fetch ↓ConsumerNetworkClient ↓Broker
Heartbeat Thread ↓Group CoordinatorClassic 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 的流水线

典型流程:
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 recordsmax.poll.records = 500可能表现为:
poll #1 → 500poll #2 → 500...而不是 Broker 必须发送六次请求。
这对慢消费调优非常重要:
- 减少
max.poll.records可以缩短单次业务处理时间; - 但不一定减少网络 Fetch 或客户端内存规模;
- Fetch 内存仍受
fetch.max.bytes、max.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 1P1 Leader → Broker 2P2 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?

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, 100Position 通常推进到:
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 = 150LEO 包含 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 有两种主要消费模式。

4.1 subscribe():Kafka 管理 Ownership
consumer.subscribe(List.of("order-events"));应用声明:
我希望消费这些 TopicKafka 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 Assignmentassign → Application Assignment4.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 fulfillmentGroup riskGroup analyticsGroup notification每个 Group 拥有自己的:
Partition AssignmentCommitted OffsetLagProcessing Lifecycle
5.2 同一 Group 内的核心约束
对传统 Consumer Group:
一个 Partition 在任一时刻只分配给同一 Group 中的一个 Consumer。
例如:
Topic = 8 PartitionsGroup = 3 Consumers可能分配:
Consumer A → P0 P1 P2Consumer B → P3 P4 P5Consumer C → P6 P7这保证:
- Partition 内 Record 由单一 Consumer 顺序拉取;
- Position 不会由同组多个实例竞争推进;
- 故障时可以把 Ownership 转移给其他成员。
5.3 Consumer 数大于 Partition 数
Partitions = 3Consumers = 5最多只有三个 Consumer 获得有效 Partition。
其余实例空闲。
因此:
有效并行 Consumer 数<=Partition 数这也是为什么看到 Lag 高时,不能无脑增加 Consumer。
5.4 一个 Consumer 可以拥有多个 Partition
Kafka 不要求一 Consumer 一 Partition。
Partitions = 12Consumers = 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 Coordinator6.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:
BrokerTopicPartitionReplicaLeaderGroup Coordinator
管理消费组状态:
MemberSubscriptionAssignmentOffset CommitGroup Epoch一个属于集群控制面,一个属于消费组协调。
6.4 Coordinator 内存状态与持久状态
Coordinator 会缓存:
Group MetadataCommitted 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 自动分配”。
它是一套明确的协议。

7.1 FindCoordinator
Consumer 先向任意 Broker 查询:
group.id 对应的 Coordinator 在哪里获得:
Coordinator Node后续 Group 请求发给该 Node。
如果 Coordinator 发生变化,Client 会标记旧 Coordinator 不可用并重新查找。
7.2 JoinGroup:汇报成员能力和订阅
每个 Consumer 向 Coordinator 发送:
Member ID / Instance IDSubscription支持的 Assignment StrategiesProtocol MetadataRebalance TimeoutCoordinator 收集当前 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 P1Member B → P2 P3这就是 Classic Protocol 为什么把较多协调复杂性放在 Client。
7.4 SyncGroup:提交并分发 Assignment
Group Leader 通过 SyncGroup 把每个成员的 Assignment 发给 Coordinator。
其他成员也发送 SyncGroup,等待 Coordinator 返回自己的 Assignment。
成功后 Group 进入:
StableConsumer 获得 Partition,初始化 Position,开始 Fetch。
7.5 Generation 的作用
一次成功 Rebalance 会形成新的:
GenerationOffset 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 状态机

经典 Group 状态通常包括:
EmptyPreparingRebalanceCompletingRebalanceStableDead8.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,正常:
HeartbeatFetchCommit8.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 暴露:
PreparingRebalanceCompletingRebalanceEmptyStableDead等状态计数指标。5
9. Partition Assignor:分配不是简单平均除法
Kafka 必须决定:
Members × Subscriptions × Partitions ↓Assignment不同 Assignor 优化目标不同。

9.1 RangeAssignor
按 Topic 分别分配连续 Partition 区间。
例如两个 Topic:
Topic A: A0 A1 A2Topic B: B0 B1 B2Consumers: C1 C2可能:
C1 → A0 A1 B0 B1C2 → 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 列表为:
RangeAssignorCooperativeStickyAssignor这样可以通过滚动升级移除 Range,把所有 Consumer 切换到 CooperativeSticky。2
9.5 Assignor 选择必须考虑什么
不要只比较“均不均匀”。
至少评估:
Topic 数量Subscription 是否一致Partition 数量Assignment 稳定性Stateful Resource 初始化成本Partition 迁移代价Rebalance 停顿容忍度滚动升级兼容例如 Consumer 为每个 Partition 建立本地缓存或数据库连接时,Sticky 的价值会明显增加。
10. 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→ 确认释放→ Assign10.3 Cooperative 为什么可能需要多轮
假设:
C1 当前拥有 P0 P1C2 当前拥有 P2 P3新增 C3目标:
C1 → P0C2 → P2C3 → P1 P3第一轮先让:
C1 revoke P1C2 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

启用:
group.protocol=consumer当前默认仍是:
group.protocol=classic11.1 为什么重新设计协议
Classic Protocol 的关键限制包括:
- Assignment 在 Client Group Leader 上计算;
- JoinGroup / SyncGroup 形成全组同步屏障;
- 大 Group 的 Subscription Metadata 很重;
- 一个成员变化容易影响全组;
- Client 配置控制 Heartbeat、Session 和 Assignor;
- Broker 很难独立增量推进每个成员状态。
新协议目标:
更大 Consumer Group更短 RebalanceBroker 端集中 AssignmentFully Incremental简化 Client改进线程模型11.2 Broker-side Assignor
新协议中,Assignment Strategy 由 Broker 控制:
group.consumer.assignors=uniform,rangeKafka 4.3 默认提供:
uniformrange默认优先:
uniformConsumer 可以通过:
group.remote.assignor=range请求 Broker 端可用 Assignor。1
这与 Classic 的:
partition.assignment.strategy完全不同。
11.3 Member Epoch
Classic Protocol 常围绕:
Generation管理全组版本。
新协议更强调:
Member EpochTarget AssignmentCurrent 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 主要状态为:
EmptyAssigningReconcilingStableDead其中:
Assigning
Broker 正在计算或更新目标分配。
Reconciling
成员当前 Assignment 正在向目标 Assignment 收敛。
这比 Classic 的 Preparing / Completing 更直接表达增量协调过程。5
11.6 Heartbeat 与 Session 配置移动到 Broker
使用新协议时,以下 Client 配置不再支持:
heartbeat.interval.mssession.timeout.mspartition.assignment.strategyenforceRebalance()改由 Broker 配置:
group.consumer.heartbeat.interval.msgroup.consumer.session.timeout.msgroup.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 最容易混淆的三个时间概念。

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→ CommitFailedException12.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-01Consumer 成为静态成员。
短暂重启时,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.msmax.poll.interval.ms改得非常大,只会:
- 延迟真实故障接管;
- 延迟 Partition 释放;
- 扩大 Lag;
- 隐藏业务线程卡死;
- 让发布和扩缩容变慢。
正确做法是先测量:
Poll 间隔分布Batch Processing P99Heartbeat RTTRebalance 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 移除成员 ↓Rebalance13.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.strategygroup.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 latencyRebalance durationGroup convergence业务处理停顿不要在 Callback 内启动不受控并发后立刻返回,否则可能在旧 Ownership 工作仍未结束时让新 Owner 开始处理。
15. Offset Commit:什么时候才算“消费完成”

15.1 Commit 的对象是恢复位置
假设成功处理:
P0 Offset 0..99应提交:
P0 → 100含义:
重启后从 100 开始。
Commit 并不会删除 Kafka 中的 Record。
Retention 仍由 Topic 配置决定。
15.2 Offset Commit 如何持久化
流程:
Consumer ↓ OffsetCommitRequestGroup Coordinator ↓ Append__consumer_offsets ↓ ReplicationCommitResponse官方实现文档说明,Coordinator 把 Offset Commit 追加到 __consumer_offsets Compacted Topic,并在达到内部 Topic 的复制确认条件后返回成功;Coordinator 同时维护 Offset Cache 加速查询。4
15.3 enable.auto.commit=true
Kafka 4.3 默认:
enable.auto.commit=trueauto.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 常见组合
业务循环中:
周期性 commitAsyncRebalance / Close 时:
最后一次 commitSync这样兼顾吞吐和退出确定性。
但仍要检查当前 Assignment 与处理水位。
15.9 多线程处理时提交什么
假设 Poll Thread 把同一 Partition 的 Record 交给 Worker:
Offset 100 完成Offset 101 仍处理中Offset 102 完成不能提交:
103否则 101 失败后无法重放。
可提交的水位是:
从旧 Position 开始,已经连续完成的最大前缀的下一 Offset。
即:
Completed: 100, 102Gap: 101Safe Commit = 101这就是多线程 Consumer 最难的部分之一。
15.10 Commit 不能解决外部副作用 Exactly Once
Kafka Record ↓更新 MySQL ↓Commit Kafka OffsetMySQL 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.reset16.2 auto.offset.reset
Kafka 4.3 支持:
earliestlatestby_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 Offset16.5 Replay 的两种方式
在线 Consumer 内 seek
适合:
- 小范围调试;
- 自定义恢复;
- Rebalance Listener 中外部 Offset 初始化。
管理工具 Reset Group Offset
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 慢频繁 Rebalance17.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关键规则:
- 同一 Partition 的业务顺序必须受控;
- Poll Thread 仍需持续 poll;
- Worker Queue 满时 pause;
- 只提交连续完成前缀;
- Revoke 时停止接收新任务并等待 / 取消旧任务;
- 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-ratefetch-latency-avg / maxfetch-size-avg / maxrecords-per-request-avgbytes-consumed-raterecords-consumed-rate判断:
- Fetch 是否频繁但 Batch 很小;
- Broker 响应是否变慢;
- 网络是否成为瓶颈;
- Consumer 是否真正有数据流入。
18.2 Lag 层
records-lagrecords-lag-maxrecords-lead必须按 Partition 观察。
总 Lag 可能掩盖:
一个 Hot Partition+大量空闲 Partition18.3 Poll 与处理层
last-poll-seconds-agotime-between-poll-avg / maxpoll-idle-ratio-avgrecords-consumed-rate如果:
time-between-poll-max接近 max.poll.interval.ms说明业务处理已经接近失去 Ownership。
18.4 Coordinator 与 Rebalance
Classic / Consumer 指标关注:
assigned-partitionsheartbeat-rateheartbeat-response-time-maxlast-heartbeat-seconds-agorebalance-totalrebalance-rate-per-hourlast-rebalance-seconds-agofailed-rebalance-totalBroker 侧 Kafka 4.3 还暴露:
consumer-group-count{state}consumer-group-rebalance-rateclassic group completed rebalance rateoffset-commit-rate等 Group Coordinator 指标。5
18.5 Commit
commit-ratecommit-latency-avgcommit-latency-maxcommit-failed-rateCommit 慢可能意味着:
- 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 故障诊断树

19.1 症状:Lag 持续上升
按顺序检查:
Ownership
Group 是否 Stable?Rebalance Rate 是否异常?Assigned Partitions 是否减少?是否存在空闲 Consumer?Fetch
Leader 是否可用?Fetch Latency 是否升高?是否有大 Batch 卡住?Broker Quota / TLS 是否异常?Processing
DB / HTTP P99Worker QueueGCCPUHot PartitionCommit
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=consumer20.1 公共入口
阅读:
KafkaConsumerConsumerConsumerDelegateConsumerConfigSubscriptionStateConsumerMetadata重点问题:
- KafkaConsumer 如何选择 Delegate?
- subscribe、assign、poll、commit、seek 如何委托?
- SubscriptionState 保存哪些 Position 与 Assignment?
- Thread Ownership 如何检查?
- wakeup 如何跨线程生效?
20.2 Classic Client 路线
建议阅读:
ClassicKafkaConsumerConsumerCoordinatorAbstractCoordinatorHeartbeatHeartbeatThreadConsumerNetworkClientFetcher / FetchCollectorSubscriptionState关注:
poll()
updateAssignmentMetadataIfNeededpollForFetchesprocessBackgroundEvents / callbacksreturn ConsumerRecords具体方法名会随版本演进,应固定 Kafka 4.3 Tag。
ConsumerCoordinator
ensureCoordinatorReadyensureActiveGrouponJoinPrepareperformAssignmentonJoinCompletecommitOffsetsSync / AsyncmaybeAutoCommitAbstractCoordinator
FindCoordinatorJoinGroupSyncGroupHeartbeat ThreadGeneration / Member IDRejoin Need20.3 新 Async Consumer 路线
阅读:
AsyncKafkaConsumerApplicationEventHandlerConsumerNetworkThreadNetworkClientDelegateConsumerMembershipManagerHeartbeatRequestManagerCommitRequestManagerFetchRequestManagerOffsetsRequestManagerBackgroundEventProcessor关注:
- Application API 如何转换为 Event?
- ConsumerNetworkThread 如何 poll 网络?
- Heartbeat Response 如何推进 Member Epoch?
- Target Assignment 与 Current Assignment 在哪里 Reconcile?
- Rebalance Callback 如何回到应用线程?
- Commit Future 如何跨线程完成?
- Fetch Buffer 如何交付给 poll?
20.4 Broker Group Coordinator 路线
阅读:
GroupCoordinatorServiceGroupCoordinatorShardGroupMetadataManagerOffsetMetadataManagerClassicGroupConsumerGroupConsumerGroupMemberTargetAssignmentBuilder关注:
Classic Group
PreparingRebalanceCompletingRebalanceStableGenerationProtocolLeaderConsumer Group
Group EpochMember EpochTarget AssignmentCurrent AssignmentAssigning / Reconciling / StableOffset
OffsetCommit__consumer_offsets RecordOffset CacheOffset FetchExpiration / Tombstone20.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 数万行代码。
先建立:
PositionMembershipAssignmentCommit四个状态域。
20.6 源码阅读记录模板
每条链固定记录:
入口 API ↓运行线程 ↓输入状态 ↓发出的 Request / Event ↓Broker 状态变化 ↓Client 本地状态变化 ↓回调 / 返回值 ↓错误和重试边界例如 Rebalance:
入口:poll()线程:Application + Heartbeat / Network输入:Subscription / Member Epoch请求:JoinGroup 或 ConsumerGroupHeartbeatBroker:Target AssignmentClient:Revoke / Assign / Position Init回调:ConsumerRebalanceListener错误:Coordinator Change / Fencing / Timeout21. 可运行实验:亲手观察 Consumer Group 与 Rebalance
下面使用单节点 Kafka 4.3.0 做本地实验。
环境定位:学习和协议观察,不代表生产部署建议。
21.1 目录结构
kafka-consumer-lab/├── docker-compose.yml├── pom.xml└── src/main/java/lab/ConsumerLab.java21.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:启动:
docker compose up -ddocker logs -f kafka-ch0421.3 创建六 Partition Topic
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create \ --topic order-events \ --partitions 6 \ --replication-factor 1验证:
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --describe \ --topic order-events21.4 发送带 Key 的测试消息
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
打开三个终端。
终端一:
mvn -q exec:java \ -Dexec.args="--instance c1 --group group-a --protocol classic"终端二:
mvn -q exec:java \ -Dexec.args="--instance c2 --group group-a --protocol classic"终端三:
mvn -q exec:java \ -Dexec.args="--instance c3 --group group-a --protocol classic"观察:
六个 Partition 如何分给三个 Consumer新增 Consumer 时哪些 Partition 被 RevokeCooperativeSticky 是否保留已有 Assignment每个 Partition 是否只出现在一个 Member 的 Assignment 中停止其中一个 Consumer,再观察其 Partition 如何被其他成员接管。
21.8 实验二:不同 Group 独立消费
另开终端:
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:
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe \ --group group-a重点查看:
TOPICPARTITIONCURRENT-OFFSETLOG-END-OFFSETLAGCONSUMER-IDHOSTCLIENT-ID查看成员和 Assignment:
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe \ --group group-a \ --members \ --verbose查看状态:
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe \ --group group-a \ --state21.10 实验四:制造慢消费与 max.poll.interval.ms 超时
启动慢 Consumer:
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
启动第一个成员:
mvn -q exec:java \ -Dexec.args="--instance n1 --group new-group --protocol consumer"再启动:
mvn -q exec:java \ -Dexec.args="--instance n2 --group new-group --protocol consumer"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 是否仍能继续处理。
查看状态:
docker exec kafka-ch04 \ /opt/kafka/bin/kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --describe \ --group new-group \ --state21.12 实验六:Offset Reset 与 Replay
先停止 group-a 全部 Consumer。
Dry Run:
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确认后执行:
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 实验七:比较自动提交与手动提交
自动提交:
mvn -q exec:java \ -Dexec.args="--instance auto-1 --group auto-group --protocol classic --manual-commit false --sleep-ms 1000"在处理过程中强制终止进程,再重启。
手动提交:
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 具有:
并行 FetchCompleted Fetch本地缓存分批 poll 返回一次 poll 可能完全消费本地缓存,也可能只推进网络和协调状态而不返回 Record。
22.2 “Heartbeat 正常就不会 Rebalance”
错误。
Heartbeat 证明:
成员进程和网络仍能维持 Session但应用长时间不 poll 仍会触发:
max.poll.interval.msKafka 需要同时识别 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=consumer22.6 “Consumer 越多,吞吐一定越高”
错误。
传统 Consumer Group 的有效并行度受:
Partition 数Hot PartitionBroker 能力下游能力共同限制。
22.7 “把 Timeout 调大就能解决 Rebalance”
通常只是延迟真实故障发现。
Timeout 调整必须建立在:
Poll Interval DistributionProcessing P99Heartbeat RTTRebalance Duration这些测量结果上。
22.8 “手动 Commit 就不会重复消费”
错误。
业务成功后、Commit 成功前发生故障,仍会重新消费。
Manual Commit≠Exactly Once Business Effect22.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 ↓CommitFailedExceptionRecord 会被新 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→ Stable24.2 Offset 验收
能够明确区分:
Record OffsetConsumer PositionCommitted OffsetLog End OffsetHigh WatermarkLast Stable Offset并说明:
- 哪些是 Consumer 本地状态;
- 哪些是持久恢复状态;
- 哪些是 Broker Log 边界;
read_committed使用哪个可见上界。
24.3 Group 验收
能够解释:
Group IDMember IDStatic Instance IDGenerationMember EpochSubscriptionAssignmentPartition OwnershipGroup Coordinator以及这些概念分别属于 Classic 还是新 Consumer Protocol。
24.4 API 验收
能够解释:
subscribe()assign()poll()pause()resume()seek()position()committed()commitSync()commitAsync()wakeup()close()每个 API 改变了哪一类状态。
24.5 配置验收
可以把以下配置放回对应机制,而不是背参数表:
group.idgroup.protocolgroup.instance.idpartition.assignment.strategygroup.remote.assignorheartbeat.interval.mssession.timeout.msmax.poll.interval.msmax.poll.recordsenable.auto.commitauto.commit.interval.msauto.offset.resetfetch.min.bytesfetch.max.wait.msfetch.max.bytesmax.partition.fetch.bytesisolation.level24.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→ ConsumerNetworkClientConsumer Protocol
KafkaConsumer→ AsyncKafkaConsumer→ ApplicationEventHandler→ ConsumerNetworkThread→ Membership / Heartbeat / Commit / Fetch ManagersBroker 侧能够说出:
GroupCoordinatorService→ GroupCoordinatorShard→ GroupMetadataManager / OffsetMetadataManager→ ClassicGroup 或 ConsumerGroup24.10 实验验收
实际操作并解释:
- 同 Group 启动三个 Consumer;
- 停止一个成员并观察接管;
- 不同 Group 独立消费;
- 使用命令查看 Member、Assignment 与 Lag;
- 制造
max.poll.interval.ms超时; - 比较 Classic 与 Consumer Protocol;
- Reset Offset 并 Replay;
- 比较自动与手动 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 Processing25.2 Position 层
Record Offset ↓Consumer Position ↓Processed Watermark ↓Committed Offset ↓Failure Recovery25.3 Membership 层
Group ID ↓Group Coordinator ↓Member / Heartbeat ↓Session Liveness25.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 已交付 ↓业务真正成功 ↓什么时候安全 CommitKafka 可以管理 Partition Ownership 和 Offset Recovery。
它不能自动知道:
数据库更新是否成功HTTP 调用是否产生副作用扣款是否真正完成外部系统是否幂等下一篇,我们会进入 Kafka 最核心的 Failure Model:
Producer ACKLeader / FollowerReplicationISRHigh WatermarkLeader FailureDuplicateIdempotent ProducerTransactionExactly Once Boundary回答:
分布式系统必然失败时,Kafka 如何定义“消息已经成功”,又如何在丢失、重复、乱序和跨 Partition 原子性之间建立可推导的可靠性模型?
参考资料
- Apache Kafka 4.3 — Consumer Rebalance Protocol
- Apache Kafka 4.3 — Consumer and Share Consumer Configs
- Apache Kafka 4.3 — KafkaConsumer Javadoc
- Apache Kafka 4.3 — Distribution / Consumer Offset Tracking
- Apache Kafka 4.3 — Monitoring
- Apache Kafka 4.3 — Design
- Apache Kafka 4.3 — CooperativeStickyAssignor Javadoc
- Apache Kafka 4.3 — ConsumerPartitionAssignor Javadoc
- Apache Kafka 4.3 — Basic Kafka Operations
- Apache Kafka 4.3.0 Release Announcement
- Apache Kafka Source Repository
- Apache Kafka 4.3 — Upgrade Guide / Share Groups