Theory · Kafka · Consumer Internals
组内独占分区 · JoinGroup/SyncGroup 时序 · Rebalance 协议演进 · 位移提交与消费语义 —— 消费侧机制深挖
GroupCoordinator 定位、__consumer_offsets 50 分区、位移即消费状态(一个整数)
EAGER(stop-the-world)→ Cooperative 增量(KIP-429)→ 服务端分配新协议(KIP-848,4.0 GA)
自动提交的丢失/重复窗口推演、手动提交组合、at-most/at-least/exactly-once 实现组合
Why Consumer Groups Exist
具体场景:topic order-created 有 6 个分区,订单服务部署了 3 个实例。现在要让这 3 台机器"一起把订单处理完,每条恰好处理一次"——先看三个看起来能行的朴素方案为什么都不行。
相当于三个互不认识的读者各读一份完整日志 → 同一条订单被处理 3 次,重复扣款。这就是"广播":它适合多个不同业务各看一份,不适合同一个业务分摊工作。
每条消息都先抢锁或查一遍去重表再处理。结果:消费从"顺序读日志"退化成"随机读 + 随机写 + 锁竞争",吞吐量塌一个数量级;而且锁/数据库本身成了新的瓶颈和单点。
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 的答案是反过来——消息永不被消费删除,只靠"游标"往前推,谁来推、推到哪、换人怎么接,交给组协议。
先立住"组"的定位(第 4 页)→ 再看加入组的完整时序与协调器(第 6-7 页)→ 然后是我们最关心的两件事:Rebalance 为什么慢、怎么治理(第 8-11 页)与位移怎么提交才不丢不重(第 12-14 页)。读完你应该能回答"消费侧为什么这么设计"以及"线上消费异常该怎么排查"。
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 | 最新写入位移 减 已提交位移;衡量"消费落后多少",是消费侧最重要的监控指标 |
先读这两篇再回来,本 deck 默认你已经知道"Kafka 是一个分区化的只追加日志":
Kafka 底层原理与高吞吐 → topic、分区、存储与零拷贝全景
副本机制与数据一致性 → ISR / 副本:位移为什么也可靠
把消费想成:一个游标在一串只追加的日志上往前推。
· 谁来推 → 分配(组 + Assignor)
· 推到哪记下来 → 位移提交(Commit)
· 推的人换了怎么办 → Rebalance(协调器 + 心跳)
后面所有名词、三代协议、一堆参数,全都是这三件事在不同规模与不同故障场景下的工程化。
Rebalance 是"换人的代价"。代价比直觉大,是因为旧协议要先全员停手再重分——所以后面三代协议的演进方向只有一条主线:把"必须停手"的范围缩到最小(EAGER 全停 → KIP-429 只停受影响分区 → KIP-848 交给 broker 增量调和)。
Consumer Group
上一页那 3 个实例,只要配上同一个 groupId,6 个分区就会被自动摊给它们(每人 2 个):既不会重复处理,也不需要任何外部锁。
每个分区同一时刻只归一个消费者独占消费——天然负载均衡:消费者多则摊薄、消费者少则一人多分区。消费进度就是每分区一个已提交位移("just one number for each partition"),不需要逐条 ACK。
每个组独立维护自己的位移,互不干扰——同一份日志可被多个业务各自消费(一条数据,N 个下游视角)。这也是"日志即数据"设计的红利:消费不删数据,只推进各自的游标。
| 追问 | 答案 |
|---|---|
| 消费者数 > 分区数会怎样? | 多出来的消费者空转(不分配任何分区)仍占成员名额;消费并行度上限 = 分区数。扩容前先看分区数 |
| 为什么需要"组"? | 两件事:①横向扩展——加机器即加消费者;②位移集中管理——组代次(generation)+ 协调器统一分配,故障时重新分配有据可依 |
| 一个分区两个消费者会怎样? | 正常不会发生:分配由组协议保证独占;若绕过协议(两个不同组或手动 assign)则重复消费——assign 模式不参与组协调 |
Group Protocols
| 协议 | 分配执行方 | Rebalance 方式 | 状态(4.3 口径) |
|---|---|---|---|
| classic + EAGER | leader consumer 在客户端算分配方案 | 先全员撤销再分配:stop-the-world | 默认协议(group.protocol=classic),仍是大多数客户端的默认路径 |
| classic + Cooperative(KIP-429,2.4+) | leader consumer(客户端) | 增量协同:只迁移变动分区,无全局 STW | 可用,需全组成员均支持该 assignor |
| 新协议(KIP-848) | broker 端 assignor(服务端分配) | 增量协调:ConsumerGroupHeartbeat 单 API 闭环,几乎无 STW | 4.0 起 GA,消费端设 group.protocol=consumer 显式启用;默认仍是 classic |
把"成员管理 + 分区分配"从客户端协议(JoinGroup/SyncGroup 两跳、leader 消费者算方案)挪到 broker 的 GroupCoordinator:成员只需一个 ConsumerGroupHeartbeat 请求,协调器返回增量分配。大规模组(成千上万成员)rebalance 从"全员抖动"变成"单成员增量调和"——官方口径可到 20 倍级提速且消除 STW。
开启:group.protocol=consumer。classic 专属配置被禁用:partition.assignment.strategy、session.timeout.ms、heartbeat.interval.ms 不再支持(心跳与会话由 broker 管理)。失效清单不含 enable.auto.commit 与手动提交 API——位移提交语义不变(ConsumerConfig.java 4.3)。
JoinGroup / SyncGroup
GroupCoordinator · __consumer_offsets
__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 接任协调器,组重新加载恢复(短暂不可用,不丢位移) |
Rebalance · Triggers & Cost
| 触发源 | 具体条件 | 典型事故 |
|---|---|---|
| 成员变化 | 新成员 JoinGroup / 主动离组(close) / session.timeout.ms 内无心跳被踢 / poll 间隔超 max.poll.interval.ms 自杀离组 | 处理逻辑偶发慢 → 触发 rebalance → 分区重分配 → 更慢 → 雪崩 |
| 订阅变化 | 任一成员订阅的 topic 集合/正则匹配结果变化(leader member 比对发现) | 动态正则订阅随集群拓扑频繁触发 |
| 分区数变化 | 订阅的 topic 扩分区(元数据变化) | 扩分区本应增量,EAGER 下也是全量重分配 |
① 所有成员先撤销全部分区(无论是否受影响)→ 全组停止消费;② JoinGroup 等最慢成员到齐;③ leader 重算全量分配;④ SyncGroup 后才恢复。停顿时长 ≈ 成员数 × 协调轮次 × 最慢成员;期间位移提交堵塞,消费 lag 集中爆发。
rebalance → 分区迁移 → 新成员冷启动(建连接/取位移/预热)处理变慢 → poll 超限或心跳延迟 → 再次触发 rebalance。治理思路:把三超时调出余量、用增量协议减少迁移量、静态成员避免重启抖动(下页起展开)。
KIP-429 · Incremental Cooperative
Assignors
| Assignor | 分配算法 | 执行方 | 均衡性 / 倾斜 | 迁移量与适用 |
|---|---|---|---|---|
| RangeAssignor(默认首位) | 按 topic 逐个:分区排序后按成员顺序切连续区间 | leader consumer | 单 topic 均衡;多 topic 累积倾斜(排前成员每 topic 多拿一个) | rebalance 迁移较大;与老版本兼容性最好 |
| RoundRobinAssignor | 全部 topic 的分区统一轮转派发 | leader consumer | 均衡好;要求组内订阅一致,否则分配错乱 | 迁移大;现较少选用 |
| StickyAssignor | 均衡优先 + 尽量保留原分配 | leader consumer | 均衡且稳定 | 迁移小,但 rebalance 仍是 EAGER 全局 STW |
| CooperativeStickyAssignor | sticky + 增量协同(两轮) | leader consumer | 均衡且稳定 | 无全局 STW、迁移最小——classic 族首选 |
| KIP-848 服务端 assignor | broker 端配置的 assignor(如 range / cooperative-sticky) | GroupCoordinator(服务端) | 由服务端算法保证 | 心跳闭环增量调和;大规模组首选(4.0+ GA) |
10 分区 3 消费者:每 topic 切 4/3/3。订阅 5 个 topic 时,排第一的成员每 topic 都多拿 1 个 → 多拿 5 个分区。消费能力差异被"成员顺序"放大——多 topic 场景换 RoundRobin/Sticky 或调整成员顺序。
单 topic 少成员:Range 无妨;多 topic 要均衡:Sticky;滚动发布频繁、组大:CooperativeSticky(全组统一配);上了 4.x 且客户端全部支持:评估 KIP-848 服务端分配,rebalance 体验最好。
Storms & Mitigation
broker 侧判定:45s 内没收到心跳即踢出。heartbeat.interval.ms=3000 约为 session 的 1/15(文档建议 ≤1/3),保证抖动下仍有多次重试机会。心跳由后台线程发(不依赖 poll),但 classic 下 poll 长阻塞会拖慢响应。
客户端侧判定:两次 poll 间隔超 5 分钟 → 主动离组再触发 rebalance。处理慢的典型死法:单批 500 条处理不完 → 超时离组 → 分区重分配 → 新成员更慢 → 风暴。
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 会话过期才被踢(被踢=又一批分区迁移) |
Offsets · Commit Semantics
每 auto.commit.interval.ms=5000 提交一次,触发点在 poll() 里(拉取前检查到期,提交的是上一次 poll 返回那批的位移)。两条时间线:
| 时刻 | 推演(处理 3s/批) |
|---|---|
| poll 返回批 N | 位移仍指向批 N-1(未提交) |
| 处理中 3s | 若此时崩溃 → 重启从批 N-1 重消费 → 重复 |
| 下次 poll 前 5s 到点 | 提交批 N 的位移(此时其实已处理完) |
| 若处理 > 5s 且崩溃 | 位移已越过未处理完的记录 → 丢失窗口 |
结论:自动提交 = 默认"至少一次偏重复",批处理时间越接近/超过 5s,重复与丢失窗口越大
commitSync:阻塞等 broker 确认,可靠但拖吞吐;commitAsync:异步即返,失败只记日志可能漏提交。标准组合:处理循环内 async + 关闭/收尾时 sync 兜底。
try { while (running) { records = consumer.poll(...); process(records); // 先处理后提交 consumer.commitAsync(); // 常规路径:异步 } } finally { consumer.commitSync(); // 退出兜底:同步 consumer.close(); // 触发 LeaveGroup }
Rebalance 回调里对"被撤销分区" commitSync(onPartitionsRevoked),避免已处理未提交被新主人重消费
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 从现存最早开始),越过的数据不可恢复——长停滞场景先把位移处理清楚再恢复消费。
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) |
Pull · Poll Loop
纯轮询会忙等打爆 broker:fetch.min.bytes=1 + fetch.max.wait.ms=500——无数据时 broker 挂住请求至多 500ms,凑够 min.bytes 立即返回。追实时时 min.bytes 保持 1;吞吐优先可调大攒批。
底层 fetch 拉回的数据先进客户端缓冲:fetch.max.bytes=50MB(单请求总上限)、max.partition.fetch.bytes=1MB(单分区上限);max.poll.records=500 只控制单次 poll 返回给业务的条数——不影响网络层 fetch。
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;与批量处理耗时一起权衡 |
Cheat Sheet · 1/2
| 并行度上限 | = 分区数;消费者数超过分区数的成员空转(不消费但仍要心跳) |
| 组 → 协调器 | abs(groupId.hashCode()) % 50 定位分区 → 该分区 leader broker 即 GroupCoordinator |
| 位移存哪 | __consumer_offsets:50 分区 · 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 |
| 成员变化 | 新成员加入 / 主动离组 / 心跳超时被踢 / poll 间隔超限离组 —— 最常出事故 |
| 订阅变化 | 任一成员订阅的 topic 集合或正则匹配结果变化 |
| 分区数变化 | 订阅的 topic 扩分区(EAGER 下也是全量重分配) |
| Range(默认首位) | 单 topic 均衡;多 topic 累积倾斜(排前的成员每 topic 多拿 1 个) |
| RoundRobin | 全局轮转,均衡好;要求组内订阅一致 |
| Sticky | 均衡 + 保留原分配,迁移小;但仍是 EAGER 全局 STW |
| CooperativeSticky | sticky + 两轮增量:无全局 STW、迁移最小——classic 族首选 |
max.poll.interval.ms(或调小 max.poll.records)group.instance.id 静态成员:滚动重启不触发 rebalanceclose() 优雅离组,别等 45s 会话超时被踢
Cheat Sheet · 2/2
| 先处理后提交 | at-least-once(最常用):崩溃重消费 → 重复交给下游幂等吸收 |
| 先提交后处理 | at-most-once:崩溃后位移已越过未处理记录 → 丢失 |
| 自动提交(默认 5s) | 触发点在 poll() 里、提交的是"上一批"位移 → 时间粒度模糊,生产建议手动提交 |
| 手动提交组合 | 循环内 commitAsync + 收尾 commitSync + onPartitionsRevoked 里 sync |
| 端到端 exactly-once | 事务生产 + sendOffsetsToTransaction 把位移并入事务 + 消费端 read_committed;或 at-least-once + 全链路幂等 |
| latest(默认) | 跳到日志末端;丢历史不丢实时(新组"消费不到存量"的工单元凶) |
| earliest | 从现存最早开始;回放友好,但存量大会造成启动洪峰 |
| none | 直接抛 OffsetOutOfRange,强制人工介入 |
| session.timeout.ms = 45s | broker 判:无心跳即踢出;心跳间隔 3s(约 1/15,留足重试机会) |
| max.poll.interval.ms = 5min | 客户端判:两次 poll 间隔超限主动离组;治"处理太慢",调大它或调小 max.poll.records |
| group.initial.rebalance.delay.ms = 3s | 新组首轮 JoinGroup 的集批窗口,把启动瞬间涌入的成员合并进一轮 |
| 网络层 | fetch.min.bytes=1 + fetch.max.wait.ms=500(长轮询防空转);fetch.max.bytes=50MB · max.partition.fetch.bytes=1MB |
| 应用层 | max.poll.records=500 只截单次 poll 返回条数,不影响网络层 fetch |
onPartitionsRevoked 里 commitSync + 下游幂等poll() 仍必须调用,只是不返回数据
Interview QA · 1/2
先盖住答案自己答一遍,再展开对照——想不起来比看得顺眼记得牢。
多余成员分不到分区、空占名额(不消费但仍需心跳维持会话)。消费并行度上限=分区数;扩消费者前先确认分区余量,必要时扩分区(注意 key 映射破坏,见 producer deck)。
三类:成员变化(加入、主动退出、session 心跳超时被踢、poll 间隔超限离组);订阅变化(topic 集合或正则结果变化);订阅 topic 扩分区。前一类最常出事故——两个超时配置是治理重点。
JoinGroup:成员汇聚、协调器选首个成员为 leader 并下发全成员清单;SyncGroup:leader 上交分配方案、协调器分发各成员结果。classic 下分配在 leader 消费者客户端计算;KIP-848 起可由 broker 端 assignor 计算。
abs(groupId.hashCode()) % 50 定位内部主题分区,该分区 leader broker 即该组 GroupCoordinator。__consumer_offsets 默认 50 分区、RF=3、compact 清理,存位移与组元数据——消费状态复用 Kafka 复制体系。
EAGER:先撤销全部分区、全员停止消费、重算全量分配——STW 时长随组规模放大。Cooperative(KIP-429):revoke 与 assign 分两轮,只迁移变动分区、其余继续消费,迁移量由 sticky 最小化。两者分配都还在客户端。
session(45s)由 broker 按心跳判定,超时被踢;max.poll.interval(5min)由客户端判定两次 poll 间隔,超时主动离组。前者治"进程假死",后者治"处理太慢";处理慢要调大 interval 或减 max.poll.records,别只调心跳。
分配上收 broker:成员只需 ConsumerGroupHeartbeat 增量协调,无全局 STW,大组 rebalance 提速一个量级。4.0 起 GA 可用于生产;需显式 group.protocol=consumer,默认仍是 classic;classic 的 session/heartbeat/assignor 配置在新协议下被禁用。
①单批最坏耗时 < max.poll.interval(减 records 或调大间隔);②KIP-345 静态成员:group.instance.id 让滚动重启不触发 rebalance;③CooperativeSticky/KIP-848 减少迁移与停顿;④close() 优雅离组,避免等 45s 被踢。
Interview QA · 2/2
同样先自答:这页的题要"结论 + 一句代价/边界"才完整。
自动提交每 5s 在 poll 内触发,提交的是"上一批"位移:处理 <5s 崩溃 → 位移未推进 → 重复;处理 >5s 崩溃 → 位移已越界 → 丢失。它是模糊语义,生产链路建议手动提交。
处理循环内 commitAsync(不阻塞、失败仅记录);finally/关闭时 commitSync 兜底确保最终提交成功;onPartitionsRevoked 里对被撤销分区 commitSync。异步提交失败可能漏序——用单调递增校验或在关键节点 sync。
位移有效时永远从提交位移继续;越界/无位移时按 reset:latest 跳末端(默认)、earliest 从现存最早、none 抛异常人工介入。常见越界原因:消费停滞超保留期日志被删。reset 治理"从哪开始",不治理提交窗口。
at-least-once=先处理完再手动提交,崩溃重消费,重复靠下游唯一键/去重表吸收。先提交后处理意味着位移已越过未处理的记录,崩溃后从下一条开始——中间那部分永远跳过,即丢失。
分区被重分配后新成员从"上次提交位移"开始 → 已处理未提交的部分被重复消费。对策:onPartitionsRevoked 里 commitSync 兜底提交;Cooperative 协议只撤销迁移分区,缩小受影响面。
底层 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 行为。
下游积压/配额限流时对部分分区暂停拉取。坑:paused 期间 poll() 仍必须继续调用(维持活跃性),只是不返回数据——连 poll 一起停会被判定失活触发 rebalance。恢复用 resume。
生产侧事务(跨分区原子+fencing)+ 消费位移经 sendOffsetsToTransaction 并入同一事务 + 下游 read_committed 过滤 aborted——Kafka Streams EOS 即此实现。不用事务的等价方案:at-least-once + 全链路幂等(唯一键/版本号)。
Related & References
组定位 → 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/#consumerconfigs | Consumer 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/#semantics | at-most/at-least/exactly-once 官方语义定义与事务组合 |