一句话回答
Kafka 很难保证消息永远只投递一次,工程上通常采用至少一次投递,再通过业务唯一键、数据库唯一约束、条件更新或幂等记录表,保证重复消息不会重复产生业务效果。
面试考察点
- 是否理解“重复收到消息”和“业务重复执行”不是同一个概念。
- 能否解释业务成功但位点提交失败为何会重复。
- 是否知道先查询再处理存在并发窗口。
- 能否针对创建、更新、扣减、外部调用设计不同幂等方案。
- 是否理解 Kafka exactly-once 的适用边界。
重复消费是怎么产生的
1. 业务成功,位点提交失败
消费者已经写入数据库,但在提交位点前崩溃,或者提交请求超时。重启后会从旧位点重新拉取消息。
2. 消费处理时间过长
处理超过 max.poll.interval.ms,协调器认为消费者失效并触发再均衡。分区被交给新消费者后,原消费者可能仍在完成旧任务,造成并发重复处理。
3. 生产者重试
网络超时后生产者重试,若幂等生产能力未正确启用,Topic 中可能出现内容相同的两条消息。
4. 人工重置位点
故障恢复、数据修复或业务重放时回退位点,历史消息会再次消费。这是有意重复,业务系统仍需安全处理。
5. 重试 Topic 与原消息并行
设计不当时,原消费迟迟未结束,但消息已被发送到重试 Topic,两个处理链路可能同时执行。
幂等键怎么设计
幂等键必须唯一表示同一次业务意图,例如:
- 支付请求号。
- 订单号 + 操作类型。
- 库存流水号。
- 上游事件 ID。
topic + partition + offset。
优先使用业务事件 ID,因为消息迁移 Topic、重新生产或跨集群同步后,Partition 和 Offset 可能变化。同一次业务重试必须复用同一个幂等键,不能每次重新生成 UUID。
方案一:数据库唯一约束
创建类操作最适合唯一约束:
CREATE TABLE payment_result (
id BIGINT PRIMARY KEY,
request_id VARCHAR(64) NOT NULL,
order_id BIGINT NOT NULL,
status TINYINT NOT NULL,
UNIQUE KEY uk_request_id (request_id)
);
消费者重复插入时触发唯一键冲突。捕获冲突后查询已有结果,确认确实是同一业务请求,再按成功处理。
“先 SELECT 判断不存在,再 INSERT”不能代替唯一约束。两个消费者可能同时查询到不存在,随后同时执行插入。
方案二:消费记录与业务写入同事务
建立消费记录表,并把记录消息 ID 与业务变更放在同一个数据库事务:
BEGIN
INSERT 消费记录(消息 ID 唯一)
UPDATE 业务数据
COMMIT
重复消息插入消费记录失败,整个事务不再执行。消费记录不能先在事务外写入,否则进程在“记录成功、业务未执行”之间崩溃,会把未完成消息误判为已处理。
记录表要规划数据量和清理周期。清理后若旧消息仍可能重放,幂等保护会失效,因此保留时间要覆盖最大重放窗口。
方案三:状态机条件更新
订单状态推进可以利用合法前置状态:
UPDATE orders
SET status = 'PAID', paid_at = NOW(), version = version + 1
WHERE id = ?
AND status = 'PENDING'
AND version = ?;
第一次更新成功,重复消息因状态不再是 PENDING 而影响 0 行。此时要查询当前状态,区分“已经成功”“非法状态”和“订单不存在”。
状态机还能防止乱序消息让状态倒退,例如已退款订单不能再次变成已支付。
方案四:流水账而不是直接覆盖
余额和库存扣减不能只执行:
UPDATE account SET balance = balance - 100 WHERE id = ?;
重复消息会重复扣款。应先用业务流水号建立唯一流水,再根据流水更新余额,且二者在同一事务内。账务系统还需要借贷平衡、对账和冲正机制。
方案五:Redis 快速拦截
SET idempotent:{messageId} 1 NX EX 86400
Redis 可以降低重复请求进入数据库的压力,但不能作为核心业务唯一防线:
- Key 可能过期后再次消费。
- Redis 故障或数据淘汰会丢失记录。
- 设置 Key 成功后、业务提交前崩溃会造成误判。
核心数据仍应使用数据库唯一约束或原子状态更新。Redis 更适合前置过滤。
调用第三方接口怎么幂等
向支付、短信或其他服务发请求时,应把业务幂等键传给下游。如果下游支持幂等接口,重复请求返回第一次结果。
若下游不支持幂等,需要本地维护调用状态,并处理请求超时后的“不确定结果”。超时不能直接再次执行不可逆操作,应先查询下游状态或进入人工补偿。
Kafka exactly-once 能解决什么
Kafka 事务与幂等生产者可以保证 Kafka 内部“消费源 Topic、写目标 Topic、提交源位点”的原子性。Kafka Streams 可以借此提供处理语义。
但如果消费者还要写 MySQL、调用支付接口或发送短信,这些外部副作用不在 Kafka 事务中,仍需业务幂等。
批量消费的幂等处理
批量消息中可能只有一条重复或失败。可以逐条使用幂等事务,记录每个分区连续成功位点。不要因为批次最后一条成功就提交最大位点,也不要因为一条重复就回滚所有已完成业务却又推进位点。
幂等结果应该返回什么
重复请求不是系统异常。若第一次已经成功,重复消息通常应返回相同或等价的成功结果,并允许位点继续推进。若同一幂等键携带不同业务参数,则应告警并拒绝,避免把两次不同操作错误合并。
线上排查流程
- 确认消息的业务 ID、Topic、Partition、Offset。
- 查看是否发生过消费者重启、超时或再均衡。
- 检查生产端是否重复产生了不同 Offset 的消息。
- 查询幂等记录和业务流水是否在同一事务完成。
- 检查幂等 Key 是否过期或清理。
- 检查下游调用是否复用了幂等键。
- 确认重复的是消息投递,还是业务副作用。
常见误区
- 自动提交关闭并不能消除重复,只是让提交时机可控。
- 消费成功后立即提交仍可能在响应丢失时重复。
- Redis 分布式锁不等于永久幂等。
topic + partition + offset无法识别跨 Topic 重新生产的同一业务事件。- 捕获唯一键异常后直接忽略不安全,还要校验已有记录是否属于同一请求。
核心考点清单
- Kafka 可靠消费通常允许重复投递,业务必须做到结果幂等。
- 幂等键应代表业务意图,并在重试和重放时保持不变。
- 唯一约束比“先查再写”可靠,因为它能原子解决并发竞争。
- 幂等记录必须与业务写入放在同一事务。
- 状态机适合状态推进,唯一流水适合金额和库存变更。
- Kafka exactly-once 不能覆盖普通数据库和外部服务副作用。
高频追问与参考回答
追问 1:如何做到真正不重复消费?
无法承诺消费者永远只收到一次,但可以保证业务只生效一次。通过稳定幂等键、数据库唯一约束、事务和状态机,让重复投递返回已有结果。
追问 2:用 Redis 记录消息 ID 可以吗?
可以作为快速过滤,但存在过期、淘汰、故障以及与数据库事务不一致的问题,核心业务仍要有数据库约束兜底。
追问 3:为什么消费记录和业务操作必须同事务?
若分开提交,任一顺序都存在崩溃窗口:先记消费记录会漏处理,先执行业务会在重试时重复处理。同事务才能一起成功或一起回滚。
追问 4:幂等表会不会越来越大?
会,需要按业务重放周期分区或归档。清理时间必须大于消息保留、死信恢复和人工重放的最大窗口。
追问 5:重复消息和乱序消息有什么区别?
重复是同一业务事件多次到达,乱序是不同版本的事件顺序颠倒。状态机除了幂等,还需携带版本号或事件时间防止旧事件覆盖新状态。
机制全景图
下面把「Kafka 怎么保证不重复消费?」从输入到结果压缩成一条可复述的主链路。面试时先用图建立全局坐标,再进入局部实现,能避免只背零散结论。
flowchart LR
A["拉取消息"]
A --> B["执行业务副作用"]
B --> C["记录幂等键/状态"]
C --> D["提交 offset"]
D --> E["失败后重新投递并识别重复"]
完整链路:从输入到结果
沿着「拉取消息 → 执行业务副作用 → 记录幂等键/状态 → 提交 offset → 失败后重新投递并识别重复」观察输入、状态与输出,下面每个阶段都对应一个可以在源码、日志或系统表中验证的位置。
1. 拉取消息
Kafka 默认至少一次场景下,崩溃、超时和再均衡都可能让同一 offset 再次交付。
2. 执行业务副作用
业务处理可能写数据库、调用外部接口或发新消息,不同副作用需要各自的幂等证据。
3. 记录幂等键/状态
幂等键通常来自业务请求 ID,并由唯一约束或状态机在最终写入点原子裁决。
4. 提交 offset
只有业务事务成功后才提交 offset;提交成功与否仍可能结果未知,因此不能靠提交消灭重复。
5. 失败后重新投递并识别重复
重复到达时查询已有结果并返回成功,而不是再次执行;过期清理需覆盖 Kafka 最大重放窗口。
源码与实现定位
| 入口 | 阅读重点 |
|---|---|
| ConsumerCoordinator | offset 提交 |
| 数据库唯一业务键 | 最终副作用裁决 |
源码或系统表应按上表顺序追踪:先确认入口实际走到哪条路径,再用运行时数据验证,而不是仅凭类名或配置推测。
参数配置与可复现实验
UNIQUE(event_id, event_type)
在业务提交后、offset 提交前 kill 消费者,验证重放只返回旧结果。
验证步骤与预期结果
1. 固定输入和基线
先在没有故障注入的环境执行上述配置,固定数据规模、并发度、运行时版本和预热时间。以「幂等冲突」为主基线,记录值应满足「记录分区基线」;同时保存 唯一冲突/幂等命中率、处理到提交耗时,使后续变化能够回到同一时间轴比较。
2. 从实现入口确认路径
在「ConsumerCoordinator」确认请求确实进入「offset 提交」对应的实现,再沿「数据库唯一业务键」观察「最终副作用裁决」。如果入口路径都未命中,就不应继续调整下游参数,而应先检查调用条件、版本或路由是否与假设一致。
3. 注入本文特有的失败模式
优先复现「用内存 Set 去重重启后失效」,并把单一变量逐级放大,直到「幂等冲突」越过「超过基线 2 倍」。随后再分别验证「幂等记录和业务写不在同一事务」和「幂等键 TTL 小于消息回放周期」,三类故障分开执行,避免多个变量同时变化而无法归因。
4. 执行止损和根因修复
第一轮只应用「最终写唯一约束」,确认它能控制影响范围;第二轮应用「记录处理中/成功状态」,验证核心链路恢复;最后落实「保留期覆盖最大回放」,消除同类问题再次出现的条件。每一步都保留变更前后数据,不用“感觉变快了”替代测量。
5. 通过退出条件
实验只有同时满足三项才算通过:「幂等冲突」回到「记录分区基线」、「P99 延迟」回到「小于业务预算」、「端到端差异」回到「0」,并且业务结果差异为零。若性能恢复但结果不一致,仍应视为失败;若指标恢复后很快再次越线,则说明只完成了临时止损,没有消除根因。
量化基线
| 指标 | 样例基线/口径 | 风险线 | 结论 |
|---|---|---|---|
| 幂等冲突 | 记录分区基线 | 超过基线 2 倍 | 重复投递 |
| P99 延迟 | 小于业务预算 | 突破预算 | 重复投递 |
| 端到端差异 | 0 | 任意非零 | 停止并对账 |
这些数值是实验口径或示例告警线,不是可复制到所有系统的固定答案;上线阈值应由本系统稳态、峰值和故障演练共同确定。
事故复盘:优惠券发放偶发重复
消费者先调用券系统再记录已处理标记,进程在两步之间宕机,重试再次发券。将券请求号作为下游幂等键,并在本地数据库用唯一约束记录处理状态后,跨系统两端都有防线。
| 失败模式 | 首要证据 | 第一处置动作 |
|---|---|---|
| 用内存 Set 去重重启后失效 | 唯一冲突/幂等命中率 | 最终写唯一约束 |
| 幂等记录和业务写不在同一事务 | 处理到提交耗时 | 记录处理中/成功状态 |
| 幂等键 TTL 小于消息回放周期 | 重复投递来源 | 保留期覆盖最大回放 |
发布与回滚检查点
- 发布前:确认「ConsumerCoordinator」对应实现和上述配置在目标版本仍然有效,并保存「幂等冲突」基线。
- 灰度中:同时观察 唯一冲突/幂等命中率、处理到提交耗时、重复投递来源;任一指标越过表中风险线,就停止继续扩量。
- 回滚时:先执行「最终写唯一约束」控制影响,再回退代码或参数;涉及持久状态时必须额外核对结果差异。
- 发布后:至少覆盖一个完整峰值周期,确认「用内存 Set 去重重启后失效」没有再次出现,才关闭变更观察窗口。
方案对比与选型
| 方案 | 更适合的场景 | 主要收益 | 代价与边界 |
|---|---|---|---|
| 数据库唯一约束 | 副作用最终落同一数据库 | 正确性证据强、实现直接 | 冲突处理与热点索引成本 |
| 幂等状态表 | 多步骤业务需记录处理中/成功 | 状态可追踪、可补偿 | 表增长与清理治理 |
| 下游幂等 API | 副作用在外部系统 | 真正保护最终资源 | 依赖对方契约和保留窗口 |
选型至少带上 消息速率、峰值带宽、分区数、消息大小和积压恢复时间,并用上面的量化基线验证;未知数据应明确为待测假设。
设计边界与工程取舍
幂等不是“忽略重复”,而是相同业务请求无论执行多少次都收敛到同一结果,并能返回原结果或可解释状态。
工程落地遵循:可靠性来自生产、Broker、消费和业务幂等的完整闭环。回答时直接引用「ConsumerCoordinator」、配置实验和事故数据,比复述固定模板更有说服力。