尧图精选

事务Outbox模式+CronJob兜底:彻底解决微服务消息丢失问题

🕒 发布时间:2026/9/19 6:37:36 📁 来源:尧图网络
先聊聊一个我上个月刚排查完的生产事故。线上订单服务一切正常用户下单也提示成功了可积分服务那边死活没收到订单创建消息。对账一跑用户积分没加上客服工单瞬间堆起来。我拉日志一层层往下查最后发现消息在事务提交成功 发送动作执行之间悄无声息地没了。这种问题在微服务架构里太典型了它不是偶发是只要你用写库发消息这种组合逻辑就一定会踩到的坑。这一篇是系列第106篇我把自己在几个项目里沉淀下来的方案完整拆一遍事务性 Outbox 模式配合兜底 CronJob。Outbox 负责从根上解决业务数据和消息不一致的问题CronJob 则用来捞那些因为进程崩溃、网络抖动、MQ 故障而滞留在表里的漏网之鱼。全文不绕弯子直接讲表结构、代码实现、调度配置和我在生产环境里踩过的坑适合所有被消息丢失折磨过的后端开发也适合正在设计异步消息链路的架构师参考。1. 消息丢失问题的根源1.1 一个典型的事务与消息矛盾场景很多团队最初写下单逻辑代码大概长这样Transactional public void createOrder(OrderCommand cmd) { Order order new Order(cmd); orderMapper.insert(order); mqTemplate.send(ORDER_TOPIC, order.getId()); }这个写法看似没啥毛病但仔细拆一下就有两个大坑。第一个坑如果mqTemplate.send在事务提交之前执行消息已经发出去了但数据库事务还没提交消费者如果响应够快立刻拿这个消息里的订单号去查订单表极大概率查不到数据。第二个坑如果事务提交之后才执行send但发送这一步抛了异常按照Transactional的默认回滚规则订单数据会一起回滚用户下单直接失败。但只要加try-catch把异常吞掉或者把发送动作放到afterCommit里数据提交成功了消息却丢了下游永远不知道有新订单产生。有人可能会说把发送放到TransactionSynchronizationManager.registerSynchronization的afterCommit回调里不是已经挺好了吗方向对了但还不够。afterCommit只保证发送动作在事务提交后触发它不保证发送一定成功。MQ 宕机、网络超时、应用进程在回调前崩溃消息照样丢掉。说白了这个方案把发送当成了一次普通调用失败了没有持久化残留也没有重试的依据。1.2 消息丢失的常见路径在生产环境里我见过太多种消息丢失的路径远不止代码顺序写错这一种发送时网络超时客户端抛异常但服务端其实已经写入成功。此时如果你选择不重试消息就看起来丢了如果你选择重试又容易造成重复。MQ 集群开启异步刷盘消息写入 PageCache 就返回成功节点宕机时未落盘的消息全部丢失。这个属于运维配置问题但业务方经常不感知。业务代码里直接把 MQ 发送写在事务内事务回滚了消息却已经发出去下游消费后产生脏数据对账时感觉像消息错乱。消费者处理失败但异常被吞掉没有抛出也没有开启手动 ack消息被确认后丢弃。运维清理队列时误操作或者磁盘写满后 broker 自动丢弃消息。这些路径里有一个很扎心的共性绝大多数团队从来不会为发消息这个动作设计持久化和重试机制。大家默认 MQ 是可靠的结果就是事故发生时连个排查的抓手都没有。1.3 强一致不可行最终一致性才是答案聊到这里其实已经触及了分布式系统最核心的问题本地事务和远程消息发送是两回事它们之间没有共同的原子性边界。想要让数据库里的订单数据和 MQ 里的消息做到强一致传统做法是引 XA 分布式事务。但 XA 协议对数据库和 MQ 都有严格要求锁粒度大、性能损耗高生产环境里几乎没有团队愿意为发条消息付出这么大的代价。所以业内的主流方案退一步选择最终一致性——先想办法确保消息一定不丢再通过重试让它在有限时间内送达。Outbox 模式就是在这个背景下流行起来的。它的核心思路非常朴素既然远程发送不可靠那我就先把消息以数据库记录的形式固化下来再通过一个独立的异步组件把固化下来的消息搬运到 MQ。只要业务数据和 Outbox 记录在同一个本地事务里提交数据一致性的问题就从业务库里的一条记录 MQ 里的一条消息变成了业务库里的一条记录 业务库里的另一条记录——后者处于同一个事务边界内天然一致。2. 事务Outbox模式设计与落地2.1 Outbox表结构怎么设计既然要持久化消息第一件事就是建表。我见过五花八门的 Outbox 表设计但核心字段基本逃不开下面这套CREATE TABLE outbox_event ( id bigint(20) unsigned NOT NULL AUTO_INCREMENT, event_id varchar(64) NOT NULL COMMENT 全局唯一事件ID用于消费端幂等, aggregate_type varchar(64) NOT NULL COMMENT 聚合类型如 ORDER、PAYMENT, aggregate_id varchar(64) NOT NULL COMMENT 聚合根ID如订单号, event_type varchar(128) NOT NULL COMMENT 事件类型如 ORDER_CREATED, payload json NOT NULL COMMENT 事件载荷业务自定义结构, status tinyint(4) NOT NULL DEFAULT 0 COMMENT 0待发送 1已发送 2死信, retry_count int(11) NOT NULL DEFAULT 0 COMMENT 已重试次数, next_retry_time datetime NOT NULL COMMENT 下次重试时间, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_event_id (event_id), KEY idx_status_next_retry (status,next_retry_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT事务消息Outbox表;几个关键设计点我说一下。event_id必须全局唯一而且要加唯一索引。它不光是给这条记录做标识更重要的是给下游消费者做幂等用的。因为后续无论是发布器重试还是 CronJob 兜底补偿都可能把同一条消息发送多次消费者拿到event_id才能安全地去重。aggregate_type和aggregate_id用来标记这条事件属于哪个业务对象排查问题时可以根据订单号快速找到所有关联事件。payload用 JSON 类型存储是 MySQL 5.7 的推荐做法MySQL 8 里还能直接对 JSON 字段做查询方便调试。索引设计上(status, next_retry_time)这个复合索引是核心。发布器和兜底任务查询待发送记录的条件几乎都是WHERE status 0 AND next_retry_time NOW()这个复合索引能完美命中。如果你不建这个索引表数据一多扫描接口就会拖垮数据库。2.2 事务内写Outbox的关键实现表建好了关键是业务代码里怎么写入。我这里直接给出我项目中在用的写法Service public class OrderService { Resource private OrderMapper orderMapper; Resource private OutboxEventMapper outboxEventMapper; Transactional(rollbackFor Exception.class) public void createOrder(OrderCreateCommand cmd) { // 1. 业务主操作插入订单 Order order new Order(cmd); orderMapper.insert(order); // 2. 构造事件并写入outbox与订单插入在同一个事务内 OrderCreatedEvent event OrderCreatedEvent.builder() .orderId(order.getId()) .userId(order.getUserId()) .amount(order.getAmount()) .build(); OutboxEvent record OutboxEvent.builder() .eventId(event.getEventId()) .aggregateType(ORDER) .aggregateId(order.getId()) .eventType(ORDER_CREATED) .payload(JsonUtils.toJson(event)) .status(OutboxStatus.PENDING) .nextRetryTime(new Date()) .retryCount(0) .build(); outboxEventMapper.insert(record); } }这段代码是整套方案的灵魂它有两个硬性要求。第一Transactional注解必须同时包住业务插入和 Outbox 插入两者共享同一个数据库连接和事务边界。这样只要事务提交订单数据和 Outbox 记录一定同时可见事务回滚两者一起消失。第二事务内绝对不要发 MQ。因为事务未提交时消费者可能查不到数据而且 MQ 发送慢会拖长事务把数据库连接耗死。如果你的业务足够复杂一个方法里要操作多张表事务边界可以用编程式事务TransactionTemplate来控制效果一样只是代码写法更灵活。我实际项目里遇到过在同一个事务里写了两条 Outbox 记录的情况这完全没有问题只需按照聚合事件逐一插入即可。2.3 发布器实现与退避策略数据进表之后需要一个组件把 PENDING 状态的记录搬运到 MQ。这个组件我习惯叫它 Publisher用Scheduled驱动核心代码如下Component public class OutboxPublisher { private static final int BATCH_SIZE 100; Resource private OutboxEventMapper outboxEventMapper; Resource private MQProducer mqProducer; Scheduled(fixedDelay 1000) public void publishPendingEvents() { ListOutboxEvent pendingEvents outboxEventMapper.scanPending( OutboxStatus.PENDING, new Date(), BATCH_SIZE ); for (OutboxEvent event : pendingEvents) { publishSingle(event); } } private void publishSingle(OutboxEvent event) { try { mqProducer.send(event.getEventType(), event.getPayload()); outboxEventMapper.markSent(event.getId(), event.getEventId()); } catch (Exception e) { log.warn(outbox消息发送失败, eventId{}, retryCount{}, event.getEventId(), event.getRetryCount(), e); int nextRetryCount event.getRetryCount() 1; Date nextRetryTime calcNextRetryTime(nextRetryCount); outboxEventMapper.updateRetryInfo(event.getId(), nextRetryCount, nextRetryTime); } } private Date calcNextRetryTime(int retryCount) { long[] delays {1000, 5000, 30_000, 120_000, 900_000, 3_600_000, 21_600_000}; int idx Math.min(retryCount - 1, delays.length - 1); return new Date(System.currentTimeMillis() delays[idx]); } }发布器每秒扫描一次每次最多取 100 条发送成功就把状态改成 SENT发送失败就更新重试次数和下次重试时间。这里有几个细节值得展开发送成功后是按UPDATE outbox_event SET status1 WHERE id?标记还是直接删除是很多团队要做的第一个选择。我推荐的做法是标记 定期清理保留已发送记录至少 3 天方便排查问题如果直接删除发送成功后应用刚好崩溃重启后扫描看不到这条记录就无法确定它是否已经发到 MQ对账只能靠下游。退避策略用的是简单的递增数组不要引入复杂的随机退避和高精度重试算法没必要。默认 1 秒、5 秒、30 秒、2 分钟、15 分钟、1 小时、6 小时重试到上限后进入死信状态。这个策略在多次故障演练中表现稳定既能尽快速恢复又不会因为大量重试把 MQ 流量打满。3. CronJob兜底机制的落地3.1 为什么发布器之外还需要兜底发布器每秒都在扫表听起来已经足够可靠了。但现实是我把 Outbox 方案上线半年后还是遇到过一次消息丢失事故。那次是发布器所在的应用实例发生了 OOM整个 JVM 卡死Scheduled的线程池当然也停止了执行。更麻烦的是监控面板只盯着 MQ 的发送量完全没有注意到 Outbox 表里的 PENDING 记录正在默默增长。等到对账发现问题的时候已经过了快一个小时。这个事故给了我一个教训所有在业务应用进程内运行的定时任务都会因为进程自身的死亡而失效。Scheduled跑得再快它的生命周期和宿主进程绑定在一起。所以我们需要一个独立于业务应用之外的兜底任务它必须活在不同的进程里或者由像 Kubernetes 这样的容器平台来调度从而保证即使业务应用全部挂掉兜底任务依然能按照自己的节奏运行。3.2 CronJob调度设计与代码实现兜底任务我强烈建议用 K8s CronJob 来承载而不是在业务应用中再写一个Scheduled。两个原因第一应用内定时任务在多副本部署时会导致多个实例同时执行补偿逻辑需要引入分布式锁复杂度翻倍第二K8s CronJob 由容器平台调度天然保证同一个时刻只有一个 Pod 在执行任务并且失败、重试、历史记录都看得见。CronJob 的配置如下apiVersion: batch/v1 kind: CronJob metadata: name: outbox-compensate-job spec: schedule: */5 * * * * concurrencyPolicy: Forbid startingDeadlineSeconds: 300 successfulJobsHistoryLimit: 3 failedJobsHistoryLimit: 3 jobTemplate: spec: template: spec: restartPolicy: OnFailure containers: - name: compensator image: your-registry/outbox-compensator:1.0.0 args: [--jobcompensate]concurrencyPolicy: Forbid一定要配置它的意思是上一次执行还没结束时下一次调度直接跳过防止两个补偿任务同时跑导致同一批消息被重复发送。schedule字段是 5 分钟执行一次这个频率和发布器的 1 秒扫描叠加能保证最坏情况下消息延迟在 5 分钟级别。补偿任务的核心代码Component public class OutboxCompensateJob { private static final int MAX_BATCH 500; private static final long STALE_THRESHOLD_MS 10 * 60 * 1000; Resource private OutboxEventMapper outboxEventMapper; Resource private MQProducer mqProducer; public void compensate() { Date deadline new Date(System.currentTimeMillis() - STALE_THRESHOLD_MS); ListOutboxEvent staleEvents outboxEventMapper.scanStale( OutboxStatus.PENDING, deadline, MAX_BATCH ); for (OutboxEvent event : staleEvents) { try { mqProducer.send(event.getEventType(), event.getPayload()); outboxEventMapper.markSent(event.getId(), event.getEventId()); log.info(兜底补偿发送成功, eventId{}, event.getEventId()); } catch (Exception e) { handleCompensateFail(event, e); } } } private void handleCompensateFail(OutboxEvent event, Exception e) { log.error(兜底补偿发送失败, eventId{}, retryCount{}, event.getEventId(), event.getRetryCount(), e); if (event.getRetryCount() 10) { outboxEventMapper.markDead(event.getId()); alert(outbox死信事件, event); } else { outboxEventMapper.updateRetryInfo( event.getId(), event.getRetryCount() 1, new Date(System.currentTimeMillis() 60_000) ); } } }这个兜底任务有一个非常关键的设计点它不是扫描所有 PENDING 记录而是只扫描创建时间已经超过阈值的 PENDING 记录。因为正在被发布器正常处理的消息通常在 1 秒内就会从表中消失。如果兜底任务也去抢这些新记录就会和发布器形成竞争反而造成大量重复发送。STALE_THRESHOLD_MS我设置为 10 分钟这个值能保证兜底任务只处理那些发布器搞不定的消息同时不会让故障恢复时间拖得过长。补偿任务里如果发现某条消息重试次数已经超过 10 次就直接把状态置为死信然后触发告警。告警通道可以是钉钉、企业微信、邮件一定要把event_id、aggregate_id、event_type都拼进去值班同学拿到信息就能直接去查。3.3 幂等与去重处理兜底任务和发布器本质上是同一类操作只是触发时机不同。所以同一个event_id完全可能被发送两次发布器第一次发送成功了但因为标记 SENT 的 SQL 还没来得及执行进程就挂了5 分钟后兜底任务扫描到这条仍然处于 PENDING 的记录又发送了一次。这种重复无法从发送端完全避免只能靠消费端幂等来解决。消费端幂等我推荐用数据库唯一索引 业务表字段的方式。比如积分服务要消费 ORDER_CREATED 消息给用户加积分可以在积分流水表里加一个source_event_id字段并建唯一索引整个处理逻辑放在一个事务里Transactional public void handleOrderCreated(OrderCreatedEvent event) { PointsRecord record new PointsRecord(); record.setUserId(event.getUserId()); record.setAmount(event.getAmount().intValue()); record.setSourceEventId(event.getEventId()); // 唯一索引保证同一event_id最多插入一次 pointsRecordMapper.insert(record); }这个方法的关键点在于积分流水插入成功 消费处理成功。如果同一事件重复消费第二次插入时唯一索引会直接报 DuplicateKeyException事务回滚业务不受影响。相比先查询再判断的写法唯一索引兜底没有并发窗口问题是最干净的幂等实现。用 Redis 做幂等键也可以但有一个隐患如果 setnx 成功之后业务处理过程中抛异常key 已经存在了重试会被拒掉而第一次其实没有成功。除非你用 Lua 脚本保证检查、处理、删除的原子性否则不建议把核心业务逻辑的幂等押在 Redis 上。数据库唯一索引最大的优势是它和业务操作处在同一个本地事务里要么一起成功要么一起回滚没有中间状态。4. 常见问题与排查技巧4.1 消息重复与幂等实战记录消息重复这件事我见过最典型的一幕是某天下午订单服务发布新版本滚动重启过程中正好有几十条 Outbox 记录处于 PENDING 状态。发布器在实例 A 上发送了一条消息标记 SENT 之前实例 A 开始优雅停机SQL 没执行完实例 B 启动后发布器扫描到这些陈旧的 PENDING 记录又发了一遍。下游积分服务没有幂等保护用户积分翻倍加对账系统半夜报警。这个事故告诉我们只要发布器重启、CronJob 触发、网络重试三个因素任意叠加重复就是必然的。所以排查重复消息时第一步不是去改发送端代码而是先去确认消费端到底有没有幂等保护。如果已经建了唯一索引直接看日志里的 DuplicateKeyException 就能确认重复来自哪里。如果没有幂等保护那首先要补的也是消费端而不是无休止地优化发送端。4.2 时序与延迟问题的取舍Outbox 模式天然会给消息链路引入额外延迟因为发布器是轮询扫描CronJob 是定期补偿两者都不是实时推送。我用fixedDelay 1000的发布器正常情况下消息延迟在 1 秒左右如果发布器挂掉最长可能要等 5 分钟让 CronJob 发现并补偿。这个延迟量级对大多数业务订单通知、积分变动、数据同步完全可接受但对风控、限流这类毫秒级实时场景就不适合了。时序问题更隐蔽。假设同一订单连续发生 ORDER_CREATED 和 ORDER_PAID 两个事件发布器在扫描时恰好在同一批里遍历顺序按id升序正常是先发前者再发后者。但假如 ORDER_CREATED 发送失败ORDER_PAID 反而先到了消费者消费端处理逻辑如果对顺序敏感就会出现订单还没创建却先处理了已支付的脏数据。解决思路有两种一是给事件带上occurred_at业务时间戳消费端按业务时间排序二是用aggregate_id做分区键保证同一聚合的事件在 MQ 中严格有序但这要求 MQ 支持消息分区。实际项目中我倾向于第一种简单、直观、不容易被 MQ 特性绑死。4.3 积压与死信的快速定位方法生产环境最怕的就是 Outbox 表里的 PENDING 记录暴涨。不要一上来就盯着 MQ 控制台看先在数据库里跑一条聚合 SQLSELECT status, retry_count, COUNT(*) AS cnt FROM outbox_event WHERE created_at NOW() - INTERVAL 1 HOUR GROUP BY status, retry_count;这条 SQL 能直接告诉我们积压的类型。如果大量记录retry_count 0那说明发布器根本没扫到这些记录问题大概率出在应用进程本身——可能是应用挂掉了也可能是Scheduled调度的线程池被耗尽了。如果retry_count很大那说明发布器一直在尝试但 MQ 一直不健康应该去看 broker 的磁盘、连接数、分区状态。如果死信记录status 2特别多那说明某一条链路已经彻底坏了优先去看死信记录里出现频率最高的event_type顺着它查消费端日志。我把排查过程中最常见的几种现象整理成了一张速查表现象可能原因处理方法outbox表持续膨胀已发送记录未清理增加定时清理任务保留3~7天retry_count高但MQ显示正常消息体格式问题消费端反序列化失败查看消费端异常栈检查payload字段同一条消息被重复消费发布器重试或CronJob补偿确认消费端event_id唯一索引CronJob一直不执行镜像拉取失败或K8s调度被挂起检查CronJob资源状态、Pod事件PENDING积压但retry_count低发布器死锁或调度线程池耗尽抓线程dump查看应用日志这里还要提一下清理已发送记录的正确姿势。直接执行DELETE FROM outbox_event WHERE status 1在数据量大时会锁住大量行影响业务写入。分批删才安全DELETE FROM outbox_event WHERE created_at NOW() - INTERVAL 7 DAY AND status IN (1, 2) LIMIT 1000;配合一个每天凌晨执行的定时任务就能把表大小控制在一个合理的量级。写在最后的实战体会从第一次踩到消息丢失的坑到把 Outbox 和 CronJob 这套组合拳完整落地我最大的体会是没有一套方案能解决所有问题但有持久化保障的重试机制能解决 99% 的消息丢失问题。剩下的 1%要交给消费端幂等和链路监控来兜底。还有一个小技巧值得分享。很多人只关注 Outbox 表的发送成功率却忽略了一个重要信号PENDING 数量的变化趋势。我后来专门做了一块监控面板展示 Outbox 表 PENDING 数量、最老记录的等待时间、死信数量这三个指标。PENDING 数量持续上升说明发布器或者上层进程出了问题PENDING 数量平稳但死信上升说明消费端业务逻辑出了问题。这个面板比单纯看 MQ 的发送 QPS 要灵敏得多几乎每次都能比业务方更早发现问题。如果你正在规划消息可靠性方案我的建议是先把 Outbox 表建起来事务注解确定好边界然后立刻补上独立于应用进程的 CronJob 兜底别等项目上线之后再来补作业。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →