跳到主要内容

幂等消费、顺序与重复处理

At-least-once 消息系统会在确认丢失、消费者崩溃和重新均衡等情况下重复投递。消费者要把“同一业务动作再次到达”设计成正常输入,并明确顺序只在哪个业务键范围内成立。

1. 业务完成与消息确认之间存在故障窗口

消费者处理一条扣款消息:

  1. 数据库事务提交扣款。
  2. 消费者提交 offset 或发送 ack。

若进程在两步之间崩溃,消息会再次到达,而第一次扣款已经生效。交换顺序则会在确认后、扣款前崩溃时丢失业务处理。

只要业务数据库与消息进度不在同一个原子提交边界,就需要接受“可能重复”或“可能丢失”中的一边,并通过额外协议解决。

2. 幂等键标识业务动作

消息应携带稳定的 event ID 或 command ID。生产者重试同一个业务动作时继续使用原 ID,不要每次生成新 ID。

消费者可以在业务数据库中建立唯一处理记录:

INSERT INTO consumed_event(event_id, consumed_at)
VALUES (?, NOW());

将去重记录和业务变化放进同一个本地事务。唯一键冲突表示该事件已完成,可以直接返回成功。若先写去重记录、后在另一个事务修改业务,后者失败会让消息永远被当成已处理。

去重表需要保留策略。删除太早会让历史重放再次产生副作用,永久保留则会持续增长,应根据消息最大重放窗口和审计要求分区或归档。

3. 业务状态转换可以天然幂等

有些操作不需要单独去重表,可以通过条件更新表达允许的状态变化:

UPDATE orders
SET status = 'PAID', paid_at = ?
WHERE order_id = ? AND status = 'PENDING';

第一次消息完成转换,重复消息因条件不满足而不再修改。消费者还要区分“已经是目标状态”和“出现非法逆向状态”,避免把任何受影响行数为 0 都当成成功。

金额累加、发送邮件等非幂等副作用通常需要唯一业务记录或下游幂等接口。

4. Exactly-once 需要说明作用范围

Kafka 事务可以把读取 Kafka、写入 Kafka 和提交消费 offset 放入同一 Kafka 事务,并配合 read_committed 隔离减少下游观察到的中间或重复结果。

它不自动覆盖 MySQL 提交、HTTP 调用、对象存储和邮件发送。跨出 Kafka 边界后,仍需本地事务、Outbox、幂等键或可补偿状态机。面试或设计文档提到 exactly-once 时,应明确是哪一段数据流、在哪个存储系统内成立。

5. 顺序通常按业务键保证

全局顺序会限制并行能力,多数业务只需要同一个订单、账户或设备的事件有序。

Kafka 可以把业务键稳定映射到同一 partition,并让一个 consumer group 中该 partition 由一个消费者处理。若消费者把消息再分发到多线程,需要按 key 串行或在写入端校验版本,否则完成顺序仍可能改变。

RabbitMQ 即使从单 queue 按顺序投递,多个消费者、不同处理时长和 requeue 也会改变完成顺序。严格顺序往往意味着同一键只能有一个在途消息,并要决定毒消息是否阻塞后续消息。

6. 版本号防止旧消息覆盖新状态

网络和重试会让版本 3 的事件先于版本 2 到达。消费者可以在消息中携带聚合版本,并使用条件写入:

UPDATE search_projection
SET payload = ?, version = 3
WHERE aggregate_id = ? AND version < 3;

旧版本随后到达时不会覆盖新状态。若版本必须连续,可以暂存缺口并等待前序消息;若投影只关心最新结果,可以跳过旧版本并通过定期对账修正遗漏。

时间戳不一定适合作为版本,因为多节点时钟可能偏移,两个更新也可能拥有相同精度值。优先使用由权威写入端生成的单调版本。

7. 重新均衡和超时也会产生重复

Kafka consumer 处理时间超过轮询或会话相关限制时,partition 可能分配给其他消费者;旧消费者若继续完成副作用,就会与新消费者重复处理。批处理还可能在部分记录完成后只提交旧 offset,重启后重放整个区间。

应控制单批大小和处理时间,在分区撤销时停止接收新工作,并确保未完成任务不会在失去所有权后继续无约束写入。长任务可以拆成持久化作业,由消息只负责创建任务。

8. 用结果核对发现静默问题

消费者指标应包含处理成功、幂等命中、版本跳过、非法状态、重试、死信和最老 lag。对关键链路还要周期性比较事件源与业务结果,例如检查已支付订单是否都有记账记录。

幂等命中突然升高可能说明生产者重复发送、确认超时或消费者频繁重启。把它当成正常而不监控,会掩盖系统故障。

9. 常见问题

9.1 Java 方法加 synchronized 能实现幂等消费吗

不能覆盖多实例、重启和历史重放。幂等状态必须保存在所有消费者共享且可靠的存储中,并与业务结果形成原子或可恢复的关系。

9.2 消费者只处理一条消息时是否天然有序

单线程可以保证当前进程的处理顺序,但重试、重新入队、partition 重新分配和生产端跨分区写入仍可能改变到达顺序。要从生产分区、消费并发和结果版本三处共同定义顺序。

10. 面试题

10.1 如何保证消息幂等消费,并处理重复与乱序

出现公司:京东、字节跳动

考察重点

  • 业务提交与 ack、offset 提交之间的故障窗口。
  • 去重记录、条件状态转换和本地事务。
  • partition 顺序、版本号与 exactly-once 的作用范围。

相关内容:第 1 节“业务完成与消息确认之间存在故障窗口”至第 8 节“用结果核对发现静默问题”。

参考回答

At-least-once 消费中,业务提交后、消息确认前崩溃会导致重复投递。消息应携带稳定事件 ID,消费者把去重记录和业务变化放入同一个本地事务,或者用带前置状态的条件更新让重复执行不再产生副作用。幂等状态不能只放在进程内。

顺序通常按业务键保证:Kafka 把同一键写入同一 partition,消费者内部也要按键串行;乱序仍可通过权威版本号做条件更新,防止旧事件覆盖新状态。Kafka exactly-once 可以覆盖 Kafka 内的读写和 offset,写 MySQL 或调用第三方时仍需幂等、Outbox 或补偿流程。