消息重复消费怎么办?怎么实现幂等?
一句话回答
主流消息队列提供的是至少一次投递:生产者重试、消费者处理完还没提交 offset 就崩溃、rebalance,都会让同一条消息被消费多次。所以重复不可避免,只能让消费逻辑幂等:同一条消息处理一次和处理多次,结果一样。常用手段有:业务唯一键 + 唯一索引、去重表和业务操作放在同一个事务里、状态机条件更新(UPDATE ... WHERE status = ?)、乐观锁版本号,以及用 Redis SET NX 做前置过滤。Kafka 的生产者幂等只能消除生产端重试带来的重复,管不了消费端。
详细解析
重复从哪来
| 环节 | 场景 |
|---|---|
| 生产端 | 消息已写入 Broker,但确认在网络上丢了,生产者重试又发了一次 |
| 消费端 | 业务处理完、offset 或 ack 还没提交时进程崩溃,重启后重新消费 |
| rebalance | 消费者处理太慢被踢出消费者组,分区分给别人,未提交的消息被再消费一次 |
| 业务层 | 本地消息表投递成功但标记失败,定时任务又投递了一次 |
保证消息不丢失 的手段(重试、处理完再提交)本身就会制造重复,两者是一体的:至少一次 + 幂等 = 效果上的恰好一次。
幂等方案
① 唯一键 + 唯一索引:最简单可靠。比如"订单支付成功后发放积分",积分流水表对 order_id 建唯一索引,重复消息插入时报唯一键冲突,捕获后直接当作成功。
② 去重表:业务表没有天然的唯一键时,单独建一张消费记录表,以消息 ID(或业务键)为主键。插入去重记录和业务操作必须在同一个本地事务里:业务失败时去重记录一起回滚,下次重试还能正常处理。
③ 状态机条件更新:带状态的业务,把状态检查写进 UPDATE 条件:
-- 只有"待支付"的订单才能改成"已支付",重复消息影响行数为 0
UPDATE orders SET status = 'PAID', paid_at = NOW() WHERE id = ? AND status = 'UNPAID';
根据影响行数判断:1 表示本次处理生效,0 表示已经处理过(或状态不对),直接确认消息即可。
④ 乐观锁版本号:UPDATE ... SET version = version + 1 WHERE id = ? AND version = ?。消息里带上生产时的版本号,重复消息因版本不匹配而更新失败。适合"把数据更新为某个快照"的场景。
⑤ Redis SET NX:消费前 SET msg:{id} 1 NX EX 86400,设置失败说明处理过或正在处理。它性能好,但和数据库不在一个事务里,有两个坑:
- 处理失败要删除这个 key,否则重试时会被误判为已处理,消息等于丢了
- key 设置成功后、业务数据提交前进程崩溃,key 存在但业务没做,重试时被误判为已处理。所以 Redis 只适合做前置过滤,最终还是要靠数据库约束兜底
| 方案 | 可靠性 | 适用场景 |
|---|---|---|
| 唯一索引 | 高 | 有天然业务唯一键的插入类操作 |
| 去重表 + 同一事务 | 高 | 通用,业务和去重记录在同一个库 |
| 状态机条件更新 | 高 | 订单、工单等有状态流转的业务 |
| 乐观锁版本号 | 高 | 更新类操作 |
| Redis SET NX | 中 | 高并发下的前置过滤,不能单独依赖 |
用哪个 ID 去重
优先用业务唯一键(订单号、支付流水号),而不是 MQ 生成的消息 ID:应用层重新发送(比如本地消息表补发、业务代码捕获异常后再调一次发送)时,MQ 会把它当作一条新消息,分配新的消息 ID,同一个业务事件就有了两个 ID。生产者也可以在消息里放一个自己生成的唯一 ID,例如本地消息表的主键。
Kafka 的幂等和事务管到哪
enable.idempotence=true 让 Broker 按"生产者 ID + 分区 + 序列号"去重,只解决生产者内部自动重试导致的重复写入,而且没配置 transactional.id 时,生产者重启后 ID 会变,重启前后发的重复消息去不掉;应用层自己再调一次 send 也去不掉。Kafka 事务配合消费者的 isolation.level=read_committed,能把"消费 + 处理 + 写回 Kafka"做成原子操作,但只限于数据都在 Kafka 里;处理结果写到 MySQL、调用外部接口,仍然要自己做幂等。
代码示例
去重表和业务操作在同一个事务里(Node.js + mysql2):
async function onPointsMessage(msg) {
const conn = await pool.getConnection()
try {
await conn.beginTransaction()
// consumed_message 以 (consumer, biz_key) 为主键,重复插入会报 ER_DUP_ENTRY
await conn.query('INSERT INTO consumed_message (consumer, biz_key) VALUES (?, ?)', ['points', msg.orderId])
await conn.query('UPDATE account SET points = points + ? WHERE user_id = ?', [msg.points, msg.userId])
await conn.commit()
} catch (err) {
await conn.rollback()
if (err.code === 'ER_DUP_ENTRY') return // 已经处理过,当作成功,正常确认消息
throw err // 其他错误抛出去,不确认消息,等待重试
} finally {
conn.release()
}
}
面试官可能追问
支付接口的幂等怎么做?
思路一样,只是"消息 ID"换成了请求的唯一标识。客户端进入支付页时先向服务端申请一个幂等号(或者直接用订单号),提交时带上;服务端用唯一索引保证同一个幂等号只生成一笔支付记录,重复请求直接返回第一次的结果。并发的重复请求可以先用 Redis SET NX 拦一下,返回"处理中"。调用第三方支付时,也把自己的流水号传过去,依赖对方按流水号去重。
去重表会越来越大怎么办?
重复消息一般只会在较短时间内出现(重试、rebalance),可以按时间分区或定期清理很久以前的记录。保留多久取决于 MQ 的消息保留时间和最长的重试周期,清理后再收到同一条消息的概率可以忽略时才能删。如果业务表本身有唯一键或状态字段,就不需要去重表。
消息乱序会影响幂等吗?
会。比如"订单已取消"先于"订单已支付"被处理,状态机条件更新会让支付消息失败,这正是想要的结果;但如果是"更新为最新数据"的消息,旧消息后到会覆盖新数据,这时要用版本号或更新时间做条件,只接受比当前更新的数据。顺序问题见 消息顺序与积压。
易错点
- 先查"是否处理过"再执行业务,查和改不是原子的,并发时两个消费者都会查到"未处理"
- 用 Redis 记录消息 ID 后,处理失败没有删除 key,导致消息永远不再处理
- 去重记录和业务操作不在一个事务里,业务失败了去重记录却留下了
- 以为开了 Kafka 生产者幂等,消费端就不用做幂等
AI 模拟面试官
用自己的话回答,AI 对照参考答案打分、指出遗漏,再追问,最多 3 轮
这道题你掌握了吗?
选一个最接近的状态,没掌握的题会出现在"我的进度 · 待复习"里。
学习记录暂存在本机浏览器。登录后自动同步到账号,换设备也能看到。