考察点
这题看的是你对"至少一次"语义有没有真正的工程体感。面试官期待你指出根因——位移提交和业务处理之间存在窗口,而不是泛泛说"网络问题"。追问会落在:自动提交具体怎么丢怎么重、手动提交选同步还是异步、Rebalance 时的重复、以及下游怎么做到幂等。
参考答案
根因:两件事不是原子的
Consumer 的消费位点(offset)是它自己提交的,存在 __consumer_offsets 这个内部主题里(0.9 之前在老 ZK 里,早就不用了)。消费一条消息的完整动作其实是两步:处理消息、提交位移。这两步之间没有原子性,中间任何一个环节出事,就会出现"位移和消费实况不一致",重复或丢失由此而来。
三种典型场景:
- 处理完了,位移没提交成功(提交请求失败、Consumer 崩溃):重启或再均衡后从上次提交的位移继续消费,已处理过的消息被再消费一遍——重复消费。
- 位移先提交了,处理失败:这条消息被跳过——丢消息。
- Rebalance 触发的重复:分区被重新分配给别的 Consumer 时,如果上一个 Consumer 处理了一批但位移还没来得及提交,新接手的 Consumer 会从旧位移重新拉,这一批全部重复。
所以 Kafka 只保证"至少一次","恰好一次"要靠你自己把位移提交和处理绑成原子操作,或者让下游幂等。
自动提交:方便但不可控
enable.auto.commit=true(默认值)配合 auto.commit.interval.ms(默认 5 秒),由 poll 循环在后台定期提交"上次 poll 返回的最大位移"。它的问题在于提交时机和你的处理进度完全解耦:
- poll 出来 500 条,处理到第 200 条时崩了,但后台线程刚提交过 500 条对应的位移——300 条丢了。
- 反过来,处理完一批但还没到下次自动提交就崩了——这批全重复。
自动提交还有一个隐蔽的坑:它是"poll 到就提交位移",如果处理逻辑是异步线程做的,poll 位移早就超前于真实处理进度了,丢消息丢得更狠。生产环境处理逻辑稍微重一点,就该关掉自动提交。
手动提交:同步还是异步
commitSync() 是同步提交,阻塞直到 Broker 确认,失败会重试,可靠但拖慢消费吞吐;commitAsync() 是异步提交,不阻塞、不重试(重试可能覆盖掉更新的位移),失败可能丢失这次提交——但丢提交的后果只是重复消费,不是丢消息,在"先处理后提交"的框架下是可以接受的。
工程上常见的折中:正常流程用 commitAsync 保吞吐,在 Consumer 关闭前和 Rebalance 的 onPartitionsRevoked 回调里用 commitSync 兜底,把手里处理完的位移可靠地交出去。Rebalance 回调里记得在分区被收走之前完成位移提交,否则新接手的实例必然重复消费一批。
粒度再往细走:按分区提交
一次 poll 返回多个分区的消息,commitSync() 无参版本提交的是"所有分区已消费的最大位移",一个分区处理失败会拖累整体。更精细的做法是按分区逐条或逐批提交:自己维护 <TopicPartition, OffsetAndMetadata> 的 map,哪个分区处理完一批就提交哪个分区的位移。配合手动指定 max.poll.records 控制批次大小,能把失败重试的影响面控制在单个分区内。
终极兜底:下游幂等
就算位移管理做得再好,分布式环境下"恰好一次"的窗口总有理论残留(提交位移本身也可能失败),所以严肃业务的最后一道防线是下游幂等:
- 消息体里带业务唯一键(订单号、事件 ID),消费端写库用
INSERT ... ON DUPLICATE KEY UPDATE或唯一索引冲突重试。 - 或者维护一张"已消费消息表",用数据库事务把"业务写入"和"去重标记"绑在一起。
- Redis 判重适合短窗口,别拿它当长期去重存储,key 过期策略要想清楚。
一句话:位移管理把重复概率压到最低,业务幂等把重复的代价降到零,两个都要。
可能的追问
- 自动提交 + Rebalance 时最容易出什么问题? 分区被回收时自动提交的位移可能落后于实际处理进度,新接手的 Consumer 会重复消费一批;最坏情况 poll 循环里处理还没完成位移已被后台提交,直接丢消息。
- 能不能先提交位移再处理? 能,那是"至多一次"语义,处理失败就丢消息。适合日志采集这类可容忍少量丢失的场景,交易类业务不行。
- __consumer_offsets 是个什么东西? 一个内部压缩主题,50 个分区,key 是
<group, topic, partition>,value 是位移。Group Coordinator 按 group 的 hash 分配这个主题的某个分区,由对应 Broker 承担协调职责。