Theory · Kafka · Consumer Internals

Kafka 消费者与消费者组

组内独占分区 · JoinGroup/SyncGroup 时序 · Rebalance 协议演进 · 位移提交与消费语义 —— 消费侧机制深挖

组与协调

GroupCoordinator 定位、__consumer_offsets 50 分区、位移即消费状态(一个整数)

Rebalance 三代协议

EAGER(stop-the-world)→ Cooperative 增量(KIP-429)→ 服务端分配新协议(KIP-848,4.0 GA)

位移与语义

自动提交的丢失/重复窗口推演、手动提交组合、at-most/at-least/exactly-once 实现组合

这份 deck 深挖消费侧。主线:消费者组为什么存在、加入组的完整时序、位移存在哪、Rebalance 怎么触发怎么优化、三种消费语义怎么实现。和总览 deck 的分工:总览只讲 pull 模型和 fetch 参数,这里把组协议、rebalance 三代演进、位移窗口推演这些面试高频追问一次讲透。KIP-848 的版本口径都按 4.3 源码和官方文档核过:4.0 GA,但默认协议仍是 classic。

Why Consumer Groups Exist

先看麻烦:三台机器消费同一个 topic,为什么会出事

具体场景:topic order-created 有 6 个分区,订单服务部署了 3 个实例。现在要让这 3 台机器"一起把订单处理完,每条恰好处理一次"——先看三个看起来能行的朴素方案为什么都不行。

方案 A:三个实例各读全部 6 个分区

相当于三个互不认识的读者各读一份完整日志 → 同一条订单被处理 3 次,重复扣款。这就是"广播":它适合多个不同业务各看一份,不适合同一个业务分摊工作。

方案 B:加一把分布式锁 / 数据库去重表

每条消息都先抢锁或查一遍去重表再处理。结果:消费从"顺序读日志"退化成"随机读 + 随机写 + 锁竞争",吞吐量塌一个数量级;而且锁/数据库本身成了新的瓶颈和单点

方案 C:人工分工,配置写死"谁来读哪几个分区"

A 管 0-1、B 管 2-3、C 管 4-5。看着能跑,但三件事立刻失控:① C 挂了,分区 4-5 无人处理;② 加一台机器要改配置、重启全员;③ topic 扩到 8 个分区,分工表要重排。

三个问题浮出水面:谁读哪一段(分工)?② 读到哪了(进度)?③ 有人挂了/加人了怎么办(重分工)?——这三个问题就是消费者组存在的全部理由。

消费者组给出的答案:三个问题对应三件事

① 分工 → 分区分配:组内每个分区同一时刻只归一个成员独占消费,由分配策略(Assignor)算出方案;组成员共享一份分配表,不需要任何外部协调系统。
② 进度 → 已提交位移:每个分区的消费进度只是一个整数("下一条要读哪")。它存在 Kafka 自己的内部主题里,复用日志与副本机制,不需要额外数据库。
③ 重分工 → 协调器 + Rebalance:broker 上的协调器管成员进出,成员变化时在秒级内重新分配分区并接管——挂掉的实例负责的分区自动有人接。

为什么"用数据库当队列"不是替代答案

用一张 MySQL 表当队列:每条消息都要 UPDATE ... SET status=1 抢占 → 热点行锁竞争 + 扫索引 + 事后删历史,吞吐上限低且随数据量劣化;用 Redis List 的 LPOP:消息弹出后消费者崩溃就永久丢失,且读过的消息不能重放。Kafka 的答案是反过来——消息永不被消费删除,只靠"游标"往前推,谁来推、推到哪、换人怎么接,交给组协议。

本 deck 的路线

先立住"组"的定位(第 4 页)→ 再看加入组的完整时序与协调器(第 6-7 页)→ 然后是我们最关心的两件事:Rebalance 为什么慢、怎么治理(第 8-11 页)与位移怎么提交才不丢不重(第 12-14 页)。读完你应该能回答"消费侧为什么这么设计"以及"线上消费异常该怎么排查"。

动机页:用一个"3 实例消费 6 分区"的具体场景,让三个朴素方案(各读全量=重复、加锁去重=吞吐塌、人工分工=不可容错)依次失败,逼出三个真问题(分工/进度/重分工),再对应到消费者组的三件事(分配/位移/协调器+Rebalance)。补一段"数据库队列与 Redis List 为什么不行",把"日志+游标"这个设计选择讲成动机而非结论。禁止一上来就抛 topic/partition/offset/ISR 四个名词。

Prerequisites & Glossary

先把词认全:下面每一页都会用到它们

术语一句话理解(先记住这个,细节后面展开)
消费者
Consumer
读消息的客户端程序;本 deck 里就是你的业务进程里那个 KafkaConsumer 实例
分区
Partition
一个 topic 被切成若干条有序、只追加的子日志;并行度的基本单位
位移
Offset
消息在分区里的序号(0,1,2…);"已提交位移"就是该分区的消费进度
提交
Commit
把"我处理到哪了"这个数字写回 Kafka(存进内部主题 __consumer_offsets
消费者组
Consumer Group
一组用同一个 group.id 协作消费的消费者:组内分摊分区、组间互不干扰
协调器
GroupCoordinator
某个 broker 上的一个组件,专门管理"某一个组"的成员进出、分配下发、位移保存
再均衡
Rebalance
组成员/订阅/分区数变化时,重新分配分区的一整套过程(本 deck 的绝对主角)
分配策略
Assignor
决定"哪个分区给哪个成员"的算法;classic 协议下由客户端 leader 成员计算
会话与心跳
Session / Heartbeat
成员周期性地向协调器"报活";超时未报活即被判定离组,触发 Rebalance
消费滞后
Lag
最新写入位移 减 已提交位移;衡量"消费落后多少",是消费侧最重要的监控指标

如果你还不熟 topic / 分区 / 日志模型

先读这两篇再回来,本 deck 默认你已经知道"Kafka 是一个分区化的只追加日志":
Kafka 底层原理与高吞吐 → topic、分区、存储与零拷贝全景
副本机制与数据一致性 → ISR / 副本:位移为什么也可靠

一个最小心智模型(后面所有变体都是它)

把消费想成:一个游标在一串只追加的日志上往前推。
· 谁来推 → 分配(组 + Assignor)
· 推到哪记下来 → 位移提交(Commit)
· 推的人换了怎么办 → Rebalance(协调器 + 心跳)
后面所有名词、三代协议、一堆参数,全都是这三件事在不同规模与不同故障场景下的工程化。

提前建立的一个直觉

Rebalance 是"换人的代价"。代价比直觉大,是因为旧协议要先全员停手再重分——所以后面三代协议的演进方向只有一条主线:把"必须停手"的范围缩到最小(EAGER 全停 → KIP-429 只停受影响分区 → KIP-848 交给 broker 增量调和)。

阅读提示:术语不用背,遇到忘了的回来查这一页就行;真正需要背的是第 16–17 页那两张速查表。
前置页:十个术语先定义再使用,杜绝"如你所知"式跳跃(consumer/partition/offset/commit/group/coordinator/rebalance/assignor/session/lag)。右侧给两条前置 deck 链接,并把全部机制收拢成"游标在日志上往前推"这一个动作的三个问句。最后提前建立"Rebalance 演进 = 缩小必须停手的范围"这条主线,为第 8-9 页的协议对比铺路。

Consumer Group

消费者组:组内独占分区,跨组广播

上一页那 3 个实例,只要配上同一个 groupId,6 个分区就会被自动摊给它们(每人 2 个):既不会重复处理,也不需要任何外部锁。

组内(同 groupId):点对点

每个分区同一时刻只归一个消费者独占消费——天然负载均衡:消费者多则摊薄、消费者少则一人多分区。消费进度就是每分区一个已提交位移("just one number for each partition"),不需要逐条 ACK。

跨组(不同 groupId):广播

每个组独立维护自己的位移,互不干扰——同一份日志可被多个业务各自消费(一条数据,N 个下游视角)。这也是"日志即数据"设计的红利:消费不删数据,只推进各自的游标。

追问答案
消费者数 > 分区数会怎样?多出来的消费者空转(不分配任何分区)仍占成员名额;消费并行度上限 = 分区数。扩容前先看分区数
为什么需要"组"?两件事:①横向扩展——加机器即加消费者;②位移集中管理——组代次(generation)+ 协调器统一分配,故障时重新分配有据可依
一个分区两个消费者会怎样?正常不会发生:分配由组协议保证独占;若绕过协议(两个不同组或手动 assign)则重复消费——assign 模式不参与组协调
定位一句话:"消费者组 = 分区级独占的消费单元:组内分摊(点对点)、组间广播;消费状态 = 已提交位移,协调器(GroupCoordinator)负责成员管理与分配分发。"
消费者组两句话立住:组内独占分区是点对点,跨组独立位移是广播。三个追问要会:消费者比分区多就空转,所以扩容先看分区数;组存在的意义是横向扩展加位移集中管理;手动 assign 模式不参与组协调,两个消费者 assign 同一分区会重复消费。最后强调消费状态只是一个整数位移,这是 Kafka 消费侧一切设计的起点,也是后面位移提交窗口推演的根基。

Group Protocols

三代协议:classic(EAGER/Cooperative)与新协议 KIP-848

协议分配执行方Rebalance 方式状态(4.3 口径)
classic + EAGERleader consumer 在客户端算分配方案先全员撤销再分配:stop-the-world默认协议(group.protocol=classic),仍是大多数客户端的默认路径
classic + Cooperative(KIP-429,2.4+)leader consumer(客户端)增量协同:只迁移变动分区,无全局 STW可用,需全组成员均支持该 assignor
新协议(KIP-848)broker 端 assignor(服务端分配)增量协调:ConsumerGroupHeartbeat 单 API 闭环,几乎无 STW4.0 起 GA,消费端设 group.protocol=consumer 显式启用;默认仍是 classic

KIP-848 改了什么

把"成员管理 + 分区分配"从客户端协议(JoinGroup/SyncGroup 两跳、leader 消费者算方案)挪到 broker 的 GroupCoordinator:成员只需一个 ConsumerGroupHeartbeat 请求,协调器返回增量分配。大规模组(成千上万成员)rebalance 从"全员抖动"变成"单成员增量调和"——官方口径可到 20 倍级提速且消除 STW。

4.x 使用要点(源码口径)

开启:group.protocol=consumer。classic 专属配置被禁用:partition.assignment.strategysession.timeout.msheartbeat.interval.ms 不再支持(心跳与会话由 broker 管理)。失效清单不含 enable.auto.commit 与手动提交 API——位移提交语义不变(ConsumerConfig.java 4.3)。

版本口径(必背):"KIP-848 在 Kafka 4.0 GA 可用于生产,但需客户端显式 group.protocol=consumer 启用;classic 仍是默认。两代协议的组互不兼容——混部时确认各端版本。"
这页是全局地图。三代协议用一张表收拢:EAGER 全员撤销有全局停顿;KIP-429 增量协同还是客户端算方案但只迁移变动分区;KIP-848 直接把分配挪到 broker,一个心跳接口闭环,4.0 GA。重点背版本口径:GA 但默认不开,要显式 group.protocol=consumer。4.x 使用要点是稀缺细节:三个 classic 配置被禁用(分配策略、会话、心跳间隔),而 enable.auto.commit 和手动提交 API 不在失效清单里、语义不变——这条直接对照 4.3 源码的 CONSUMER_PROTOCOL_UNSUPPORTED_CONFIGS,能体现你读过源码。

JoinGroup / SyncGroup

加入组的完整时序:FindCoordinator → JoinGroup → SyncGroup → Heartbeat

classic 协议加入消费者组时序 消费者先按 groupId 定位协调器;成员经 JoinGroup 汇聚,协调器选首个成员为 leader member 并下发全组成员清单;leader 在客户端运行分配策略,经 SyncGroup 上交方案;协调器把各自分配下发给每个成员;此后成员周期性心跳维持会话,会话过期或 poll 超时触发再均衡。 Consumer A(先到,后被选为 leader) Consumer B GroupCoordinator(broker 端) EAGER rebalance 期间:全员停止消费(stop-the-world) FindCoordinator:groupId hash → __consumer_offsets 分区 leader JoinGroup(订阅 + 支持的分配协议列表) JoinGroup(同组加入) 等成员到齐(新组首轮 group.initial.rebalance.delay.ms=3000) → 选首个成员为 leader member,generation+1 JoinGroup 响应:A=leader(全成员订阅清单)· B=follower 分配发生在客户端:leader member 运行分配策略 (Range / Sticky / CooperativeSticky…)算出方案 SyncGroup:leader member 上交分配方案 SyncGroup(等待分配结果) SyncGroup 响应:A 的分区集合 SyncGroup 响应:B 的分区集合 onPartitionsAssigned → 开始消费(EAGER 下至此全员才恢复) Heartbeat:interval=3s · session.timeout.ms=45s 内必须到达 会话过期 / poll 间隔超限(max.poll.interval.ms)/ 主动离组 → 触发下一轮 rebalance(回到 JoinGroup) Rebalance 前的位移提交:协调器要求成员先提交未提交位移再撤销分区——避免"已处理未提交"被新成员重消费
这张时序图要能脱稿讲。第一步 FindCoordinator,按 groupId 的 hash 定位到 __consumer_offsets 的某个分区,那个分区的 leader broker 就是协调器。第二步 JoinGroup,所有成员汇聚,协调器选第一个到的当 leader member,并把全组成员清单发给它。第三步是关键:分配方案是 leader 消费者在客户端算的,SyncGroup 上交,协调器分发给每个人。所以面试官问"谁做分配",classic 答客户端 leader,KIP-848 答 broker。红色区域是 EAGER 的停顿窗口,全员停止消费直到 SyncGroup 完成。之后靠心跳维持会话,三个超时任何一个超了都触发 rebalance。

GroupCoordinator · __consumer_offsets

协调器与位移主题:消费状态也是 Kafka 日志

主题形态

__consumer_offsets:默认 50 分区(offsets.topic.num.partitions)、RF=3、清理策略 compact——只保留每个 key 的最新位移,天然"快照化"。

组 → 协调器映射

abs(groupId.hashCode()) % 50 定位分区 → 该分区 leader 所在 broker 即该组的 GroupCoordinator。所有 JoinGroup/SyncGroup/Heartbeat/位移提交都发给它。

记录内容

两类 key:位移记录 <group, topic, partition> → offset/元数据/提交时间;组元数据 → 成员列表、订阅、generation。写入走正常produce 流程 → 位移的可靠性同样由 ISR 复制保证

追问答案
为什么位移存 Kafka 自己,而不是 ZooKeeper/DB?复用既有复制与持久化(高吞吐小写入)、compact 天然保留最新状态、与消费链路同集群无跨系统一致性负担——旧 ZK 时代位移在 ZK,性能与规模都受限后迁入
compact 会把历史位移删光吗?每个 key 只保留最新值——正是语义所需("当前消费到哪");offsets.retention.minutes 控制无活跃组的位移保留期,过期整 key 清理
协调器挂了怎么办?位移/组元数据在 __consumer_offsets 的 ISR 里,分区 leader 切换 → 新 leader broker 接任协调器,组重新加载恢复(短暂不可用,不丢位移)
答题一句话:"__consumer_offsets 是一个 50 分区、compact 的内部主题;位移与组元数据都写进去,组按 groupId hash 映射到分区,分区 leader 就是协调器——消费状态完全复用 Kafka 自身的日志与复制体系。"
协调器和位移主题是消费侧的基础设施。三个默认值背下来:50 分区、副本因子 3、新组首轮延迟 3 秒。映射规则是面试常问:groupId 的 hash 对 50 取模定位分区,分区 leader 就是协调器。为什么位移存 Kafka 不存别处——复用复制体系、compact 保留最新、无跨系统一致性问题,这一问能体现架构权衡能力。协调器挂了不丢位移,因为位移本身在 ISR 里,新 leader 接任即可。顺带说一句:位移的可靠性由副本机制保证,这也是和 replication deck 的连接点。

Rebalance · Triggers & Cost

Rebalance:什么时候触发,EAGER 的 stop-the-world 代价

触发源具体条件典型事故
成员变化新成员 JoinGroup / 主动离组(close) / session.timeout.ms 内无心跳被踢 / poll 间隔超 max.poll.interval.ms 自杀离组处理逻辑偶发慢 → 触发 rebalance → 分区重分配 → 更慢 → 雪崩
订阅变化任一成员订阅的 topic 集合/正则匹配结果变化(leader member 比对发现)动态正则订阅随集群拓扑频繁触发
分区数变化订阅的 topic 扩分区(元数据变化)扩分区本应增量,EAGER 下也是全量重分配

EAGER 协议的代价链

① 所有成员先撤销全部分区(无论是否受影响)→ 全组停止消费;② JoinGroup 等最慢成员到齐;③ leader 重算全量分配;④ SyncGroup 后才恢复。停顿时长 ≈ 成员数 × 协调轮次 × 最慢成员;期间位移提交堵塞,消费 lag 集中爆发。

为什么"rebalance 风暴"恶性循环

rebalance → 分区迁移 → 新成员冷启动(建连接/取位移/预热)处理变慢 → poll 超限或心跳延迟 → 再次触发 rebalance。治理思路:把三超时调出余量、用增量协议减少迁移量、静态成员避免重启抖动(下页起展开)。

面试标准句:"Rebalance 本身是必要恶:EAGER 用全量撤销换实现简单,代价是全员 stop-the-world;2.4 的 KIP-429 和 4.0 GA 的 KIP-848 分别从'客户端增量'和'服务端分配'两个方向消解这个代价。"
Rebalance 触发源归三类:成员变、订阅变、分区数变。成员变化里最容易出事故的是两个超时:心跳超时被踢,和 poll 间隔超限主动离组。EAGER 的代价链要背:先全员撤销、等最慢成员、全量重算、才能恢复,停顿时长是成员数乘以最慢成员。风暴循环是经典事故题:rebalance 导致处理变慢,变慢又触发 rebalance。答题收束在三代协议的演进方向上,引出下一页的增量协同。

KIP-429 · Incremental Cooperative

EAGER vs CooperativeSticky:一轮全量撤销 → 两轮增量协调

EAGER 与 CooperativeSticky 两种 Rebalance 协议对比 EAGER 协议在触发后撤销全部分区、全员等待重分配后恢复,形成长停顿;CooperativeSticky 协议第一轮只撤销需要迁移的分区、其余继续消费,第二轮增量下发新分配,无全局停顿且迁移量最小。 EAGER(classic 历史默认) 消费中 revoke 全部 JoinGroup / SyncGroup 全员等待 恢复 t0 触发 t1 全组恢复 · 所有成员撤销所有分区(含不受影响的) · 停顿 ≈ 成员数 × 协调轮次 × 最慢成员 · 位移提交堵塞,lag 集中爆发 · 已有分配全作废,迁移量最大 适用判断 成员少、分区少、变更不频繁时够用;成员多 / 分区多 / 滚动发布频繁的场景应切换 Cooperative 或新协议 CooperativeSticky(KIP-429,2.4+) 消费中 Round1:仅 revoke 需迁移的分区 Round2 增量 assign 后 恢复 · 未被迁移的分区全程不停消费 · 两轮 JoinGroup/SyncGroup:先增量撤销、再增量分配 · sticky 保留尽量多的原分配 → 迁移量最小 · 仍属 classic 协议族:分配仍由 leader consumer 计算 启用条件 全组成员均为支持增量的客户端,且 partition.assignment.strategy 仅指定 CooperativeStickyAssignor(混用会触发多次全量 rebalance) 4.3 默认策略列表 [RangeAssignor, CooperativeStickyAssignor]:全组一致支持 cooperative 时自动进入增量模式(ConsumerConfig.java) KIP-848 的服务端分配则把"谁算方案"也上收 broker——协议形态详见第 3 页与第 8 页
左右对照讲。EAGER 是一条长停顿带:触发即撤销全部,全员等分配,恢复时已有分配全作废。Cooperative 的关键是把撤销和分配拆成两轮:第一轮只撤销要迁移走的分区,其他人照常消费;第二轮增量下发新分配。sticky 语义保证迁移量最小。两个易错点:Cooperative 仍属 classic 协议族,方案还是客户端 leader 算的,别和 KIP-848 混淆;默认策略列表里 Range 在前,混用 assignor 会触发多余的全量 rebalance,切换时要全组统一。

Assignors

分区分配策略对比:谁分配、多均衡、迁多少

Assignor分配算法执行方均衡性 / 倾斜迁移量与适用
RangeAssignor(默认首位)按 topic 逐个:分区排序后按成员顺序切连续区间leader consumer单 topic 均衡;多 topic 累积倾斜(排前成员每 topic 多拿一个)rebalance 迁移较大;与老版本兼容性最好
RoundRobinAssignor全部 topic 的分区统一轮转派发leader consumer均衡好;要求组内订阅一致,否则分配错乱迁移大;现较少选用
StickyAssignor均衡优先 + 尽量保留原分配leader consumer均衡且稳定迁移小,但 rebalance 仍是 EAGER 全局 STW
CooperativeStickyAssignorsticky + 增量协同(两轮)leader consumer均衡且稳定无全局 STW、迁移最小——classic 族首选
KIP-848 服务端 assignorbroker 端配置的 assignor(如 range / cooperative-sticky)GroupCoordinator(服务端)由服务端算法保证心跳闭环增量调和;大规模组首选(4.0+ GA)

Range 的倾斜怎么来的

10 分区 3 消费者:每 topic 切 4/3/3。订阅 5 个 topic 时,排第一的成员每 topic 都多拿 1 个 → 多拿 5 个分区。消费能力差异被"成员顺序"放大——多 topic 场景换 RoundRobin/Sticky 或调整成员顺序。

选择建议(答题模板)

单 topic 少成员:Range 无妨;多 topic 要均衡:Sticky;滚动发布频繁、组大:CooperativeSticky(全组统一配);上了 4.x 且客户端全部支持:评估 KIP-848 服务端分配,rebalance 体验最好。

易错点:"分配策略是组级协商的——leader 按全组成员共同支持的策略执行;成员配置不一致时按交集降级,混用 cooperative 与非 cooperative assignor 会导致额外全量 rebalance。"
策略对比表五行走一遍,重点是三列:谁分配、倾斜、迁移量。Range 的倾斜机制要能现场推:十个分区三个人每 topic 切四三三,五个 topic 第一人多拿五个,这就是多 topic 倾斜的来源。Sticky 和 CooperativeSticky 的区别只在协议不在算法:都是 sticky 语义,后者增量协同没有全局停顿。最后一行 KIP-848 把分配上收到 broker,是质变。选择建议按场景给,最后提醒策略是组级协商的,全组要统一。

Storms & Mitigation

三个超时的关系,以及风暴治理清单

session.timeout.ms = 45000

broker 侧判定:45s 内没收到心跳即踢出。heartbeat.interval.ms=3000 约为 session 的 1/15(文档建议 ≤1/3),保证抖动下仍有多次重试机会。心跳由后台线程发(不依赖 poll),但 classic 下 poll 长阻塞会拖慢响应。

max.poll.interval.ms = 300000

客户端侧判定:两次 poll 间隔超 5 分钟 → 主动离组再触发 rebalance。处理慢的典型死法:单批 500 条处理不完 → 超时离组 → 分区重分配 → 新成员更慢 → 风暴。

group.initial.rebalance.delay.ms = 3000

broker 侧:新组首轮 JoinGroup 等待窗口,把"服务启动瞬间涌入的成员"合并进一轮 rebalance。开发环境可调 0 加快联调。

治理手段做法
给活跃性留余量max.poll.records 调小(如 100–200)或 max.poll.interval.ms 调大(如 10–15 分钟),确保"单批最坏处理时间 < 间隔";GC/网络抖动场景相应放大 session.timeout
静态成员(KIP-345)配置 group.instance.id:滚动重启(session 内)不触发 rebalance,分区分配原样保留——发布频繁场景的特效药
换增量协议 / 新协议CooperativeSticky 消除全局 STW;KIP-848(4.0 GA)心跳闭环增量调和,大组 rebalance 提速一个量级
优雅退出consumer.wakeup() 打断阻塞处理;close() 发 LeaveGroup 主动离组——避免等 45s 会话过期才被踢(被踢=又一批分区迁移)
三个超时的角色要分清:session 是 broker 判的,心跳断就被踢;max.poll.interval 是客户端判的,两次 poll 之间太久主动离组;initial delay 是新组首轮的集批窗口,只有 broker 端有。治理四招:给单批处理时间留余量、静态成员避免发布抖动、换增量协议、优雅退出。其中 KIP-345 静态成员是加分项:group.instance.id 配上,滚动重启不触发 rebalance。排查顺序也讲一下:先看日志谁触发,再判哪类超时,对症下药。

Offsets · Commit Semantics

位移提交:自动提交的窗口推演与手动提交组合

自动提交(enable.auto.commit=true,默认)

auto.commit.interval.ms=5000 提交一次,触发点在 poll() 里(拉取前检查到期,提交的是上一次 poll 返回那批的位移)。两条时间线:

时刻推演(处理 3s/批)
poll 返回批 N位移仍指向批 N-1(未提交)
处理中 3s若此时崩溃 → 重启从批 N-1 重消费 → 重复
下次 poll 前 5s 到点提交批 N 的位移(此时其实已处理完)
若处理 > 5s 且崩溃位移已越过未处理完的记录 → 丢失窗口

结论:自动提交 = 默认"至少一次偏重复",批处理时间越接近/超过 5s,重复与丢失窗口越大

手动提交:commitSync / commitAsync 组合

commitSync:阻塞等 broker 确认,可靠但拖吞吐;commitAsync:异步即返,失败只记日志可能漏提交。标准组合:处理循环内 async + 关闭/收尾时 sync 兜底

try {
  while (running) {
    records = consumer.poll(...);
    process(records);            // 先处理后提交
    consumer.commitAsync();      // 常规路径:异步
  }
} finally {
  consumer.commitSync();         // 退出兜底:同步
  consumer.close();              // 触发 LeaveGroup
}

Rebalance 回调里对"被撤销分区" commitSync(onPartitionsRevoked),避免已处理未提交被新主人重消费

面试一句话:"先处理后提交 → at-least-once(重复交给下游幂等);先提交后处理 → at-most-once(可能丢);自动提交两者都不严格——它是延迟触发的 at-most-5s 语义,生产链路建议手动提交 + 幂等下游。"
位移提交这页是消费语义的地基。自动提交的关键事实:触发点在 poll 里,提交的是上一批的位移。推演两条时间线:处理三秒崩溃是重复,处理超过五秒崩溃可能丢,所以自动提交既不严格 at-least 也不严格 at-most,它是五秒粒度的模糊语义。手动提交标准组合:循环内 async 保吞吐,finally 里 sync 兜底,rebalance 回调里对被撤销分区 sync。这三段代码面试手写过才敢说熟。收束一句话:先处理后提交重复,先提交后处理丢,生产用手动加幂等下游。

Out of Range · auto.offset.reset

位移越界:为什么会发生,三个取值怎么选

取值行为适用与坑
latest(默认)越界时跳到日志末端(LEO)从新消息开始默认值;丢历史不丢实时。新组上线瞬间"默认从最新开始"常造成"为什么没消费到存量数据"的工单
earliest越界时跳到日志起点(log start offset)回放/补数友好;存量大会造成启动洪峰,配合限速
none越界直接抛 OffsetOutOfRange 异常,交业务处理强一致场景:强制人工介入,禁止静默跳过

越界的四种来源

保留期清理:日志段按 retention 删除,提交位移对应的 segment 已不存在(消费停滞超过保留期再恢复,最常见);② 新组/新分区无已提交位移;③ 位移被误提交(手工/工具改错);④ seek 越界。前提:只有"找不到有效位移/位移越界"时 reset 才生效——位移有效时永远从提交位移继续。

高频追问

Q:auto.offset.reset 能防丢吗?——不能:它只决定"找不到有效位移时从哪开始",处理中崩溃(位移未提交/已提交)与它无关。Q:消费停滞一周后恢复,位移指向的数据已被清理?——按 reset 走(latest 跳末端,earliest 从现存最早开始),越过的数据不可恢复——长停滞场景先把位移处理清楚再恢复消费。

答题收束:"位移有效 → 从提交位移继续;位移无效或越界 → 按 auto.offset.reset(latest/earliest/none)。它治理的是'从哪开始',不治理'处理与提交之间的窗口'——那是提交策略的事。"
位移越界页先给结论:位移有效就从提交位移继续,reset 只在找不到有效位移或越界时生效。三个取值:latest 默认跳末端,earliest 从现存最早开始,none 直接抛异常人工介入。越界来源最常见的是消费停滞超过保留期,日志段被删,提交位移指向的数据没了。两个高频追问一定要接住:reset 不防丢,它只管从哪开始;停滞一周后恢复,越过的时间窗口数据不可恢复,要先处理位移再恢复消费。这些边界说全了就是满分答案。

Delivery Semantics

消费语义:三种交付保证的实现组合与场景推演

语义实现组合代价
at-most-once先提交后处理(poll 后立即 commitSync 再处理);或自动提交 + 短间隔崩溃 → 已提交未处理 → ;不重复
at-least-once(最常用)先处理后提交(手动 commitSync/Async);崩溃重消费重复交给下游幂等(唯一键/去重表/版本号)
exactly-once① Kafka Streams / 事务链路:消费位移并入生产事务(sendOffsetsToTransaction)+ 消费端 read_committed;② 或 at-least-once + 全下游幂等等价实现事务开销;实现复杂度最高
典型场景推演
① 先提交后处理,处理中崩溃位移已越过该记录 → 重启从下一条开始 → 该记录丢失(at-most-once 的固有风险)
② 先处理后提交,提交前崩溃位移未推进 → 重启重消费该批 → 重复处理,需下游幂等吸收(at-least-once)
③ Rebalance 时位移未提交分区被撤销并重分配给新成员;新成员从上次已提交位移开始 → 已处理未提交的部分被重消费 → 重复(故 onPartitionsRevoked 里 commitSync)
答题模板:"Kafka 自身只提供'位移提交'这个原语,语义由提交时机 + 下游能力组合出来:先提交后处理是 at-most-once,先处理后提交是 at-least-once,配合事务与 read_committed(或全链路幂等)才是端到端 exactly-once。"
语义三件套是消费侧的必考题。先给实现组合表:先提交后处理最多一次会丢,先处理后提交至少一次要下游幂等,精确一次两条路——原生事务链路或者全链路幂等等价。然后三个场景推演必须熟练:先提交后处理崩溃丢,先处理后提交崩溃重复,rebalance 时未提交位移被新成员重消费。第三场景的对策就是 onPartitionsRevoked 里提交位移,和第 10 页呼应。最后强调 Kafka 只给位移原语,语义是自己组合出来的,这个观点很能体现理解深度。

Pull · Poll Loop

拉模型与 poll 循环:fetch 参数、活跃性、pause/resume

长轮询防空转

纯轮询会忙等打爆 broker:fetch.min.bytes=1 + fetch.max.wait.ms=500——无数据时 broker 挂住请求至多 500ms,凑够 min.bytes 立即返回。追实时时 min.bytes 保持 1;吞吐优先可调大攒批。

fetch 与 poll 的两层

底层 fetch 拉回的数据先进客户端缓冲:fetch.max.bytes=50MB(单请求总上限)、max.partition.fetch.bytes=1MB(单分区上限);max.poll.records=500 只控制单次 poll 返回给业务的条数——不影响网络层 fetch。

poll 的第二职责:活跃性

classic 协议下 poll 同时是"我还活着"的信号:间隔超 max.poll.interval.ms 主动离组。长处理要么调大间隔/减 max.poll.records,要么把处理挪出 poll 循环。KIP-848 下心跳/活跃性由 broker 闭环管理。

追问答案
为什么 pull 不 push?push 由 broker 掌节奏,慢消费者被打爆;pull 让消费者按自身能力拉,天然适配批量与背压(官方 Push vs. pull 一节)
pause/resume 什么时候用?本地缓冲积压/下游限流时对部分分区暂停拉取;注意 poll() 仍要继续调用维持活跃性(paused 分区不返回数据);配额、削峰、多租户隔离常用
max.poll.records 调大就好吗?单批更大 → 吞吐更好,但单批处理时间变长 → 更容易撞 max.poll.interval.ms;与批量处理耗时一起权衡
拉模型页抓四个点。第一,长轮询两个参数:min.bytes 一字节加 500 毫秒等待,这是总览里讲过的,这里带上源码默认值。第二,fetch 和 poll 是两层:网络层单请求 50 兆、单分区 1 兆,应用层单次 poll 500 条,别混。第三,poll 在 classic 协议下兼任心跳活跃性,处理慢被离组的机制就是它,这和第 9 页的风暴治理连起来。第四,pause/resume 的坑:暂停后 poll 还是要调,只是不返回数据,很多人暂停完连 poll 都停了,直接被踢出组。

Cheat Sheet · 1/2

一页带走(上):分工与换人

① 组与定位:谁读哪一段

并行度上限= 分区数;消费者数超过分区数的成员空转(不消费但仍要心跳)
组 → 协调器abs(groupId.hashCode()) % 50 定位分区 → 该分区 leader broker 即 GroupCoordinator
位移存哪__consumer_offsets50 分区 · RF=3 · compact;进度只是一个整数,复用 Kafka 自身的日志与副本
组内 / 组间组内分摊(点对点、分区独占);组间广播(各自独立位移,互不干扰)

② 三代协议:换人的代价

classic + EAGER全员撤销全部分区再重分 → 全局 STW;方案由 leader 消费者客户端
Cooperative(KIP-429,2.4+)两轮增量:只撤销需迁移的分区,其余继续消费;仍由客户端算,需全组统一 assignor
新协议(KIP-848,4.0 GA分配上收 broker:单心跳接口增量调和,几乎无 STW;需显式 group.protocol=consumer默认仍是 classic

③ Rebalance 三类触发源

成员变化新成员加入 / 主动离组 / 心跳超时被踢 / poll 间隔超限离组 —— 最常出事故
订阅变化任一成员订阅的 topic 集合或正则匹配结果变化
分区数变化订阅的 topic 扩分区(EAGER 下也是全量重分配)

④ 四个分配策略怎么选

Range(默认首位)单 topic 均衡;多 topic 累积倾斜(排前的成员每 topic 多拿 1 个)
RoundRobin全局轮转,均衡好;要求组内订阅一致
Sticky均衡 + 保留原分配,迁移小;但仍是 EAGER 全局 STW
CooperativeStickysticky + 两轮增量:无全局 STW、迁移最小——classic 族首选

⑤ 风暴治理四招

① 单批最坏耗时 < max.poll.interval.ms(或调小 max.poll.records
group.instance.id 静态成员:滚动重启不触发 rebalance
③ 换增量协议(CooperativeSticky / KIP-848)
close() 优雅离组,别等 45s 会话超时被踢
速查上页:把"组与协调器""三代协议""触发源""分配策略""治理四招"压成一页,供回看与考前扫一眼。分配策略四行与第 10 页详表互补,这里只留选择结论。

Cheat Sheet · 2/2

一页带走(下):进度、语义与线上排查

① 位移与语义:提交时机决定一切

先处理后提交at-least-once(最常用):崩溃重消费 → 重复交给下游幂等吸收
先提交后处理at-most-once:崩溃后位移已越过未处理记录 → 丢失
自动提交(默认 5s)触发点在 poll() 里、提交的是"上一批"位移 → 时间粒度模糊,生产建议手动提交
手动提交组合循环内 commitAsync + 收尾 commitSync + onPartitionsRevoked 里 sync
端到端 exactly-once事务生产 + sendOffsetsToTransaction 把位移并入事务 + 消费端 read_committed;或 at-least-once + 全链路幂等

② 位移越界:reset 只在"没有有效位移"时生效

latest(默认)跳到日志末端;丢历史不丢实时(新组"消费不到存量"的工单元凶)
earliest从现存最早开始;回放友好,但存量大会造成启动洪峰
none直接抛 OffsetOutOfRange,强制人工介入

③ 三个超时(4.3 默认)与谁在判

session.timeout.ms = 45sbroker 判:无心跳即踢出;心跳间隔 3s(约 1/15,留足重试机会)
max.poll.interval.ms = 5min客户端判:两次 poll 间隔超限主动离组;治"处理太慢",调大它或调小 max.poll.records
group.initial.rebalance.delay.ms = 3s新组首轮 JoinGroup 的集批窗口,把启动瞬间涌入的成员合并进一轮

④ fetch / poll 两层参数

网络层fetch.min.bytes=1 + fetch.max.wait.ms=500(长轮询防空转);fetch.max.bytes=50MB · max.partition.fetch.bytes=1MB
应用层max.poll.records=500 只截单次 poll 返回条数,不影响网络层 fetch

⑤ 线上排查:按现象反查

· 重复消费:位移未提交就崩溃 / rebalance 时未提交 → onPartitionsRevoked 里 commitSync + 下游幂等
· 消息丢失:先提交后处理 / 自动提交窗口(处理 > 5s) / 越界后 reset=latest
· lag 突增:先查 rebalance 风暴(日志 LeaveGroup),再看单批处理耗时
· pause 后掉组:暂停期间 poll() 仍必须调用,只是不返回数据
速查下页:提交时机决定语义这一条主线 + 越界三值 + 三个超时的判定主体 + fetch/poll 两层参数,最后给一张"按现象反查原因"的排查清单,覆盖重复/丢失/lag 突增/pause 掉组四类高频工单。

Interview QA · 1/2

组协议与 Rebalance 8 连问

先盖住答案自己答一遍,再展开对照——想不起来比看得顺眼记得牢

1 · 消费者数超过分区数会怎样?

空转并行度上限

多余成员分不到分区、空占名额(不消费但仍需心跳维持会话)。消费并行度上限=分区数;扩消费者前先确认分区余量,必要时扩分区(注意 key 映射破坏,见 producer deck)。

2 · Rebalance 的触发条件?

成员/订阅/分区数

三类:成员变化(加入、主动退出、session 心跳超时被踢、poll 间隔超限离组);订阅变化(topic 集合或正则结果变化);订阅 topic 扩分区。前一类最常出事故——两个超时配置是治理重点。

3 · JoinGroup/SyncGroup 各做什么?谁算分配?

leader member客户端分配

JoinGroup:成员汇聚、协调器选首个成员为 leader 并下发全成员清单;SyncGroup:leader 上交分配方案、协调器分发各成员结果。classic 下分配在 leader 消费者客户端计算;KIP-848 起可由 broker 端 assignor 计算。

4 · 组和协调器怎么定位?__consumer_offsets 是什么?

hash % 50compact

abs(groupId.hashCode()) % 50 定位内部主题分区,该分区 leader broker 即该组 GroupCoordinator。__consumer_offsets 默认 50 分区、RF=3、compact 清理,存位移与组元数据——消费状态复用 Kafka 复制体系。

5 · EAGER 与 CooperativeSticky 的区别?

全量撤销两轮增量

EAGER:先撤销全部分区、全员停止消费、重算全量分配——STW 时长随组规模放大。Cooperative(KIP-429):revoke 与 assign 分两轮,只迁移变动分区、其余继续消费,迁移量由 sticky 最小化。两者分配都还在客户端。

6 · max.poll.interval.ms 与 session.timeout.ms 的区别?

客户端判broker 判

session(45s)由 broker 按心跳判定,超时被踢;max.poll.interval(5min)由客户端判定两次 poll 间隔,超时主动离组。前者治"进程假死",后者治"处理太慢";处理慢要调大 interval 或减 max.poll.records,别只调心跳。

7 · KIP-848 新协议好在哪?现在什么状态?

4.0 GA默认 classic服务端分配

分配上收 broker:成员只需 ConsumerGroupHeartbeat 增量协调,无全局 STW,大组 rebalance 提速一个量级。4.0 起 GA 可用于生产;需显式 group.protocol=consumer,默认仍是 classic;classic 的 session/heartbeat/assignor 配置在新协议下被禁用。

8 · rebalance 风暴怎么治理?

余量静态成员增量协议

①单批最坏耗时 < max.poll.interval(减 records 或调大间隔);②KIP-345 静态成员:group.instance.id 让滚动重启不触发 rebalance;③CooperativeSticky/KIP-848 减少迁移与停顿;④close() 优雅离组,避免等 45s 被踢。

上半场八题覆盖组与 rebalance。第三题是核心:谁算分配,classic 是客户端 leader member,新协议是 broker,这一题答对说明协议真懂了。第六题两个超时的判定主体不同:session 是 broker 判心跳,max.poll.interval 是客户端判处理节奏,治理手段也不同。第七题背版本口径:4.0 GA、默认 classic、三个配置被禁用。第八题治理四招按顺序说,静态成员 KIP-345 是加分项。

Interview QA · 2/2

位移与消费语义 8 连问

同样先自答:这页的题要"结论 + 一句代价/边界"才完整。

9 · 自动提交为什么会丢/重复?

5s 窗口poll 内触发

自动提交每 5s 在 poll 内触发,提交的是"上一批"位移:处理 <5s 崩溃 → 位移未推进 → 重复;处理 >5s 崩溃 → 位移已越界 → 丢失。它是模糊语义,生产链路建议手动提交。

10 · commitSync 与 commitAsync 怎么组合?

循环 async收尾 sync

处理循环内 commitAsync(不阻塞、失败仅记录);finally/关闭时 commitSync 兜底确保最终提交成功;onPartitionsRevoked 里对被撤销分区 commitSync。异步提交失败可能漏序——用单调递增校验或在关键节点 sync。

11 · 位移越界怎么办?auto.offset.reset 三值?

latest 默认earliestnone

位移有效时永远从提交位移继续;越界/无位移时按 reset:latest 跳末端(默认)、earliest 从现存最早、none 抛异常人工介入。常见越界原因:消费停滞超保留期日志被删。reset 治理"从哪开始",不治理提交窗口。

12 · at-least-once 怎么实现?先提交后处理为什么丢?

先处理后提交下游幂等

at-least-once=先处理完再手动提交,崩溃重消费,重复靠下游唯一键/去重表吸收。先提交后处理意味着位移已越过未处理的记录,崩溃后从下一条开始——中间那部分永远跳过,即丢失。

13 · Rebalance 时已处理未提交的位移会怎样?

重消费revoked 回调

分区被重分配后新成员从"上次提交位移"开始 → 已处理未提交的部分被重复消费。对策:onPartitionsRevoked 里 commitSync 兜底提交;Cooperative 协议只撤销迁移分区,缩小受影响面。

14 · fetch 参数与 poll 的关系?

两层缓冲

底层 fetch:fetch.min.bytes=1/fetch.max.wait.ms=500 长轮询,fetch.max.bytes=50MB 与 max.partition.fetch.bytes=1MB 限制单请求体积;应用层 max.poll.records=500 只截返回条数。max.poll.records 不影响网络层 fetch 行为。

15 · pause/resume 的用途与坑?

限流背压poll 不能停

下游积压/配额限流时对部分分区暂停拉取。坑:paused 期间 poll() 仍必须继续调用(维持活跃性),只是不返回数据——连 poll 一起停会被判定失活触发 rebalance。恢复用 resume。

16 · 端到端 exactly-once 怎么拼出来?

事务+read_committed位移并入事务

生产侧事务(跨分区原子+fencing)+ 消费位移经 sendOffsetsToTransaction 并入同一事务 + 下游 read_committed 过滤 aborted——Kafka Streams EOS 即此实现。不用事务的等价方案:at-least-once + 全链路幂等(唯一键/版本号)。

下半场八题聚焦位移与语义。第九题自动提交窗口推演要能口算:五秒间隔、处理三秒崩溃重复、处理八秒崩溃丢。第十题组合拳代码要能手写。第十三题 rebalance 未提交位移的对策是 revoked 回调里 sync,和上半场呼应。第十六题端到端精确一次两条路:原生事务链路三件套,或者等价的全链路幂等。答题时主动说"Kafka 只提供位移原语,语义是组合出来的",这句话能让面试官停一下。

Related & References

相关知识点与参考

本库相关 deck

答题串联 · 一图流

组定位 → FindCoordinator(hash % 50)
JoinGroup/SyncGroup → leader 客户端分配 → Heartbeat
EAGER → Cooperative(KIP-429)→ 服务端分配(KIP-848,4.0 GA)
位移 → 自动提交窗口 / 手动 sync+async → 三种语义
越界 → auto.offset.reset;风暴 → 三超时 + 静态成员

参考来源(本 deck 版本敏感结论均可溯源至下列一手材料,2026-08 核实)

kafka.apache.org/documentation/#consumerconfigsConsumer configs 默认值与语义(fetch / poll / commit / reset)
.../clients/consumer/ConsumerConfig.java(4.3 源码)默认值源码核实:session=45000 · heartbeat=3000 · max.poll.interval=300000 · 默认策略列表 · KIP-848 失效清单
kafka.apache.org/43/operations/consumer-rebalance-protocol/KIP-848 新再平衡协议:4.0 GA、group.protocol=consumer 启用、classic 仍默认
cwiki.apache.org/.../KIP-429(Incremental Rebalance Protocol)增量协同两轮设计与 CooperativeStickyAssignor
cwiki.apache.org/confluence/display/KAFKA/KIP-345(Static Membership)group.instance.id 静态成员:滚动重启免 rebalance
kafka.apache.org/documentation/#semanticsat-most/at-least/exactly-once 官方语义定义与事务组合
收尾页串联消费侧主线:组定位、时序、三代协议、位移与语义、越界与风暴治理。参考材料里最值得重读的是 4.3 的 ConsumerConfig 源码注释——新协议下哪些配置失效都写在里面,以及官方 rebalance protocol 文档的 4.0 GA 口径。复习时先过答题串联五行,再回各页看机制细节。