可靠消息投递与消费的工程设计Outbox、幂等、死信兜底学习笔记。记录在跨系统数据同步里发送端如何保证消息不丢、消费端如何保证不乱不重的一整套工程设计。不涉及具体业务只谈模式与取舍。背景两个经典难题跨系统用 MQ 做数据同步绕不开两个问题发送端消息会丢。改数据库和发 MQ是两个系统没法用一个事务保证同时成功——先改库后发 MQ发失败就丢消息先发 MQ 后改库库回滚就是假消息。消费端消息会重、会乱。MQ 基本都是at-least-once至少投递一次重复投递、乱序到达是常态消费逻辑必须能扛。下面分发送端、消费端两头讲。一、发送端Transactional Outbox本地消息表1.1 问题双写不一致// 反例无法保证两步原子tx.begin();bizMapper.update(data);// 改库tx.commit();mq.send(msg);// 发 MQ —— 这一步失败/超时库已改、消息丢了跨数据库和MQ两个资源没有廉价的分布式事务。1.2 方案把发消息变成一次数据库写核心思想不直接发 MQ而是在同一个事务里往一张outbox本地消息表插一条待发送消息记录。业务数据和消息记录绑在一个本地事务里天然原子。CREATETABLEoutbox(idBIGINTPRIMARYKEYAUTO_INCREMENT,biz_typeVARCHAR(64),-- 业务类型biz_idVARCHAR(64),-- 业务主键payloadTEXT,-- 消息体 JSONstatusTINYINT,-- 0待发送 1已发送 2死信retry_countINT,next_retry_atDATETIME,UNIQUEKEYuk_biz(biz_type,biz_id)-- 幂等同一件事只发一条);tx.begin();bizMapper.update(data);// 业务数据outboxMapper.insert(buildOutbox(msg));// 消息记录status0tx.commit();// 两者一起成功 or 一起回滚原子1.3 把消息真正发出去三件套事务后 best-effort 立即发一次降延迟TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){// 等事务 commit 成功后才发try{producer.send(payload);// 发成功outboxMapper.markSent(id);// 标记已发送}catch(Exceptione){log.warn(best-effort send failed, leave to compensator,e);// 失败不管交补偿}}});为什么用afterCommit必须等事务提交成功数据真落库才发否则事务回滚了消息却发出去 假消息。补偿 Job 兜底定时扫status0 AND next_retry_atnow的记录重发发成功markSent失败retry_count并按指数退避算下次时间超上限转死信status2 告警。幂等靠(biz_type, biz_id)唯一索引重复写入冲突即跳过。1.4 ⚠️ 一个高频坑异步 Producer 的发送成功是假的很多 MQ 的 producer 是异步缓冲的如 Kafka 的send()、某些 SDK 的addUserMessage()消息进本地缓冲就返回真正投递到 broker 由后台线程做结果通过回调通知。所以下面这种写法有坑booleanokproducer.sendBuffered(payload);// 只是入缓冲成功不是broker 收到if(ok)outboxMapper.markSent(id);// ❌ 过早还没真发出去就标已发送丢消息窗口① 入缓冲后、后台线程发出前进程崩溃→ 缓冲在内存丢outbox 却已 status1补偿不管② 异步投递最终失败回调里若不处理 → 同样丢。正确做法只有拿到投递确认broker ack / 成功回调才 markSentproducer.send(payload,(meta,ex)-{if(exnull)outboxMapper.markSent(id);// 确认后才标已发elseoutboxMapper.markFailed(id);// 失败→留给补偿重试});一句话原则markSent 的依据必须是投递确认不是入缓冲/没抛异常。二、消费端幂等 手动 checkpoint2.1 at-least-once消费必然可能重复既然 MQ 至少投递一次消费逻辑必须幂等——重复消费同一条消息结果不变。常见手段唯一键去重、状态机前置条件、以数据版本覆盖。2.2 checkpoint消费进度的书签MQKafka/类似里消息按 offset 顺序排列消费者要提交读到哪了的位点这就是checkpoint / 提交 offset提交了位点前移→ 这条算消费成功不再投递不提交位点不动→ MQ 会重投相当于重试。关键配置关掉自动提交由业务按结果手动控制checkpoint.auto.commit false # 否则处理失败也会被自动提交 → 丢消息booleanokhandle(msg);if(ok)checkpointer.checkpoint(offset);// 成功才提交// 失败不提交 → MQ 重投重试2.3 “MQ 通知 回查” 替代消息带全量数据一个很实用的模式MQ 消息只当通知带个 key真实数据消费时反查生产方HTTP/RPC为准。好处天然防乱序不管消息到达顺序回查拿到的永远是生产方当前最新态天然幂等重复消费回查结果一致覆盖本地无副作用比定时全量拉取对账更实时、更省。代价多一次回查调用依赖生产方提供查询接口。三、消费端兜底自建死信表 重放 人工3.1 为什么自建RabbitMQ / RocketMQ 有原生死信队列DLQ但有些 MQ没有平台级的消费重试计数和 DLQ。这时需要业务自建一张死信表把重试计数 退避 状态 兜底自己管起来。3.2 死信表状态机0 待重试(PENDING) → 1 重试中(REPLAYING) → 2 已恢复(RESOLVED) 0 待重试 → 3 已忽略/死信(IGNORED终态)失败累加saveOrIncrRetry——首次 INSERT重复retry_count并发用唯一键冲突捕获转 UPDATE按指数退避算next_retry_at达上限转死信终态。成功清除消费成功markResolved。重放 Job定时扫status待重试 AND next_retry_atnow重新执行消费逻辑成功markResolved失败继续saveOrIncrRetry。CAS 防并发重放UPDATE ... SET status重试中 WHERE id? AND status待重试只有抢到的那个线程影响 1 行防止多实例重复重放同一条。人工兜底admin 后台可replay手动重放/ignore标脏数据忽略/resend手动补发。// 指数退避失败越多等越久并封顶privatestaticfinalint[]BACKOFF_MIN{1,5,30,120,480,1440};// 分钟intidxMath.min(retryCount-1,BACKOFF_MIN.length-1);nextRetryAtnow.plusMinutes(BACKOFF_MIN[idx]);3.3 为什么要退避而不是固定间隔狂重试防重试风暴下游挂久了若固定 1 秒/1 分钟重试堆积的死信会每轮全量砸向还没恢复的下游把它和自己压垮雪崩。退避让挂得越久、重试越稀。放弃时机合理退避 N 次能覆盖几小时~几天给下游充分恢复窗口固定短间隔很快就耗尽、把瞬时故障误判成永久失败。四、一个设计讨论消费失败重试别让两套重试打架消费失败后怎么重试有两种常见做法容易不小心同时用上、互相打架做法机制问题A. 不 checkpoint 让 MQ 重投失败就不提交位点MQ 按自己节奏重投重投节奏不可控MQ 定若还叠加自建退避表就是两套重试并行退避被 MQ 快速重投架空还重复消费B. 立刻 checkpoint 自管退避推荐失败即写死信表(待重试)立刻 checkpoint重试全程由 Job 按退避驱动单一重试驱动、退避 100% 可控、无重复消费推荐 B。关键顺序先写表(事务提交) → 再 checkpoint万一中间崩溃没 checkpointMQ 重投一次幂等重处理不丢。教训如果你既不 checkpoint 让 MQ 重投、又建了退避重放 Job等于同一条消息被两条路同时重试退避表形同虚设。要么全交给 MQ、要么全交给 Job别混。五、贯穿全局的几个通用点幂等发送端用唯一键防重发消费端用唯一键/回查/状态覆盖防重复消费。乐观锁CAS / 条件更新UPDATE ... WHERE 旧状态/旧版本靠影响行数判断有没有被并发改过比悲观锁轻、不阻塞、不死锁。best-effort vs 可靠快速路用 best-effort失败不较劲慢速路补偿/Job保证最终完成——又快又不丢。最终一致 对账兜底所有自动重试都救不回时靠人工 admin更稳的系统还会加定期全量对账把漏掉的补齐。六、一张图收尾发送端(生产者) 业务事务 { 改数据 写 outbox(PENDING) } 原子 │ 事务后 best-effort 发一次(确认才 markSent) │ 补偿 Job 扫 PENDING 重发(退避, 超限→死信) ▼ [ MQ ] ▼ 消费端(消费者) 收到(通知) → 回查生产方拿权威数据 → 幂等覆盖本地 ├ 成功 → checkpoint(提交位点) └ 失败 → 写死信表(待重试,退避) → checkpoint → 重放 Job 按退避重试 ├ 成功 → 已恢复 └ 到上限 → 死信终态 → 告警 → 人工(重放/忽略/补发)小结目标手段发不丢Transactional Outbox事务内写消息表 事务后 best-effort 补偿 JobmarkSent 认准投递确认收不乱MQ 通知 回查生产方为权威幂等消费收不重at-least-once 前提下做幂等唯一键/状态覆盖位点不乱提手动 checkpointauto.commitfalse成功才提交失败能救自建死信表退避重试 CAS 防并发重放 重放 Job 人工兜底别雪崩指数退避 封顶重试风暴防护并发安全乐观锁条件更新 / CAS本质就一句发送端用本地消息表把 MQ 投递纳入数据库事务保证不丢消费端用幂等 手动 checkpoint 自建死信兜底保证不乱不重、失败可救。