先说结论

Kafka 很难保证消息永远只投递一次,工程上通常采用至少一次投递,再通过业务唯一键、数据库唯一约束、条件更新或幂等记录表,保证重复消息不会重复产生业务效果。

重复消费是怎么产生的

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 事务中,仍需业务幂等。

批量消费的幂等处理

批量消息中可能只有一条重复或失败。可以逐条使用幂等事务,记录每个分区连续成功位点。不要因为批次最后一条成功就提交最大位点,也不要因为一条重复就回滚所有已完成业务却又推进位点。

幂等结果应该返回什么

重复请求不是系统异常。若第一次已经成功,重复消息通常应返回相同或等价的成功结果,并允许位点继续推进。若同一幂等键携带不同业务参数,则应告警并拒绝,避免把两次不同操作错误合并。

线上排查流程

  1. 确认消息的业务 ID、Topic、Partition、Offset。
  2. 查看是否发生过消费者重启、超时或再均衡。
  3. 检查生产端是否重复产生了不同 Offset 的消息。
  4. 查询幂等记录和业务流水是否在同一事务完成。
  5. 检查幂等 Key 是否过期或清理。
  6. 检查下游调用是否复用了幂等键。
  7. 确认重复的是消息投递,还是业务副作用。

容易踩坑的地方

  • 自动提交关闭并不能消除重复,只是让提交时机可控。
  • 消费成功后立即提交仍可能在响应丢失时重复。
  • Redis 分布式锁不等于永久幂等。
  • topic + partition + offset 无法识别跨 Topic 重新生产的同一业务事件。
  • 捕获唯一键异常后直接忽略不安全,还要校验已有记录是否属于同一请求。

常见问题

追问 1:如何做到真正不重复消费?

无法承诺消费者永远只收到一次,但可以保证业务只生效一次。通过稳定幂等键、数据库唯一约束、事务和状态机,让重复投递返回已有结果。

追问 2:用 Redis 记录消息 ID 可以吗?

可以作为快速过滤,但存在过期、淘汰、故障以及与数据库事务不一致的问题,核心业务仍要有数据库约束兜底。

追问 3:为什么消费记录和业务操作必须同事务?

若分开提交,任一顺序都存在崩溃窗口:先记消费记录会漏处理,先执行业务会在重试时重复处理。同事务才能一起成功或一起回滚。

追问 4:幂等表会不会越来越大?

会,需要按业务重放周期分区或归档。清理时间必须大于消息保留、死信恢复和人工重放的最大窗口。

追问 5:重复消息和乱序消息有什么区别?

重复是同一业务事件多次到达,乱序是不同版本的事件顺序颠倒。状态机除了幂等,还需携带版本号或事件时间防止旧事件覆盖新状态。