ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

事务性Outbox模式+CronJob兜底:彻底解决消息中间件丢消息问题

2026/9/19 10:48:22 拓冰建站 浏览量
事务性Outbox模式+CronJob兜底:彻底解决消息中间件丢消息问题 做后端时间长了你会发现一个扎心的事实消息中间件并不保证消息一定不丢。无论是Kafka、RocketMQ还是RabbitMQ从生产端到消费端任何一个环节抖动一下消息可能就悄无声息地消失了。等业务方拿着订单号来问“为什么积分没到账”“为什么短信没发出去”的时候你排查半天日志最后发现是消息在发出前就丢了那一刻的心情只能用“崩溃”来形容。今天这篇文章就聊一个在电商、支付、积分这类对数据一致性要求极高的场景里非常经典的组合方案事务性Outbox模式再加一个兜底的CronJob定时任务。这套方案不是银弹但在绝大多数业务场景下它是保证“本地事务和消息发送最终一致”的可靠手段。如果你是负责订单、支付、财务这类核心链路的开发或者正被分布式事务、消息可靠性问题折磨这篇内容值得你认真看两遍。1. 先把问题看清楚消息到底是在哪一步丢的要理解Outbox模式的价值得先拆解一个最常见的业务动作。假设用户下单支付成功系统要做两件事更新订单状态为“已支付”然后往消息队列发一条“支付成功”的消息下游积分服务、短信服务、数据统计服务都等着这条消息干活。代码写出来通常长这样先更新订单表然后调用MQ的send方法最后提交事务。看起来顺理成章但这里有三个隐藏的丢消息点。第一个点消息在事务提交之前就发送成功了但事务回滚了。比如订单状态更新成功消息也发出去了结果数据库提交时超时或者冲突导致回滚。下游消费端已经收到消息开始加积分了本地一看订单还是“待支付”两边就对不上了。第二个点发消息这个动作本身失败。网络抖动、MQ broker临时不可用、连接池不够send方法抛异常但业务代码如果没做好捕获和补偿订单状态已经改了消息却没发出去下游永远感知不到这次支付。第三个点最阴险——发送消息成功但发送这个动作发生在事务提交之前而事务最终提交失败。这类问题比前面两个更隐蔽因为你的日志里什么都查不到消息在MQ里能查到send ok但业务数据回滚了对账的时候才发现数据不一致。这几个问题的共性是本地数据库事务和消息发送是两个独立的系统操作没有原子性。你没办法让数据库和消息中间件在同一个事务里一起提交或回滚。Outbox模式解决的就是这个根因。1.1 本地事务和远程通信一对天生的冤家再说透一点。为什么分布式场景下“先发消息再提交事务”和“先提交事务再发消息”都不靠谱本质上是因为本地事务只能约束数据库这一亩三分地约束不了远程的MQ调用。举个生活化的例子你去银行柜台转账柜员把转出方的账扣了然后电话通知收款方银行“给他加钱”。如果电话打完收款方银行系统挂了这笔钱就凭空消失了。柜员不可能因为对方电话没接通就让转账方的扣款回滚因为扣款已经写入数据库了。这就和“事务提交了但消息没发出去”是一模一样的困境。所以业内最朴素也最有效的思路是不直接依赖远程消息发送的成败而是把要发送的消息先落进自己数据库和业务数据放进同一个本地事务里**。业务数据提交成功消息记录也就提交成功业务数据回滚消息记录也跟着回滚。至于什么时候真的把消息发出去另起一个进程从数据库里把消息捞出来再发。这就是Outbox模式的核心思想。2. 整体设计拆解Outbox 和 CronJob 怎么配合Outbox模式的全称叫Transactional Outbox直译过来就是“事务性发件箱”。虽然名字听起来很有距离感但实现思路非常朴素。在你的业务数据库里多建一张表比如叫outbox_event这张表就是业务系统和MQ之间的缓冲层。核心逻辑分三步。第一步业务操作和写Outbox表在同一个本地事务里完成。比如修改订单状态的SQL和insert一条支付成功事件记录的SQL一起执行要么都成功要么都失败天然保证原子性。第二步一个后台任务可以是独立的应用服务也可以是你自己的业务系统里开一个定时线程池扫描这张表把状态为“待发送”的记录读出来真正地发送到MQ。第三步消息发成功之后把这条记录标记为“已发送”或者直接物理删除。CronJob在这个设计里的角色是兜底。因为第二步里的后台任务无论如何稳健都有可能遇到发送失败、程序重启、MQ宕机等异常情况总有一些消息会滞留在表里。CronJob就负责定期扫描表里所有“滞留太久”的记录重新尝试发送直到成功为止。2.1 数据表结构第一个要敲定的细节表结构是整个方案的地基我给出一个经过实战打磨的版本你可以根据自己业务的字段做裁剪。CREATE TABLE outbox_event ( id bigint(20) NOT NULL AUTO_INCREMENT COMMENT 主键, event_id varchar(64) NOT NULL COMMENT 事件全局唯一ID用于幂等, biz_type varchar(32) NOT NULL COMMENT 业务类型如 ORDER_PAID, biz_id varchar(64) NOT NULL COMMENT 业务ID如订单号, payload json NOT NULL COMMENT 消息体JSON格式, status tinyint(4) NOT NULL DEFAULT 0 COMMENT 0-待发送 1-已发送 2-发送失败, retry_count int(11) NOT NULL DEFAULT 0 COMMENT 重试次数, next_retry_time datetime DEFAULT NULL COMMENT 下次重试时间, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT 创建时间, updated_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT 更新时间, PRIMARY KEY (id), KEY idx_status_next_retry (status, next_retry_time), KEY idx_event_id (event_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT事务性Outbox事件表;几个字段的设计意图说一下。event_id一定要有而且要在业务代码里显式生成通常用UUID或者雪花ID。为什么要这个字段因为下游消费端要用它做幂等判断同一个事件你重发多少次下游只能成功处理一次。这也是所有可靠消息方案里绕不开的一个话题消费幂等。status字段标记事件生命周期retry_count和next_retry_time是给CronJob用的。你不需要把所有状态列得太复杂0、1、2这三个状态足够覆盖绝大多数场景。至于已发送的记录是保留还是删除我的建议是定期清理不要一直在表里积压否则索引会膨胀、扫描性能会下降。2.2 核心写入链路事务注解和Outbox怎么无缝衔接有了表接下来就是业务代码怎么写入。这一步的关键是写业务表和写Outbox表必须在同一个事务里。Java的Spring Boot项目里最标准的做法是用Transactional注解保证同一个方法内的多个DAO操作共享一个数据库事务。举一个下单支付的简化案例。Service public class OrderService { private final OrderMapper orderMapper; private final OutboxEventMapper outboxEventMapper; public OrderService(OrderMapper orderMapper, OutboxEventMapper outboxEventMapper) { this.orderMapper orderMapper; this.outboxEventMapper outboxEventMapper; } Transactional(rollbackFor Exception.class) public void handlePaid(OrderPaidCommand command) { // 1. 更新订单状态为已支付 orderMapper.updateStatus(command.getOrderId(), OrderStatus.PAID); // 2. 构造Outbox事件 OutboxEvent event OutboxEvent.builder() .eventId(IdGenerator.generateId()) .bizType(ORDER_PAID) .bizId(command.getOrderId()) .payload(command.getPayloadJson()) .status(OutboxStatus.PENDING.getCode()) .build(); // 3. 写入Outbox表和上面的订单更新同事务提交 outboxEventMapper.insert(event); } }这段代码最核心的一句话在方法上的Transactional注解。它保证了第一步更新订单和第二步插入Outbox记录要么一起成功要么一起失败。订单状态改了但Outbox记录没写进去的情况不可能发生Outbox记录写进去了但订单状态没改的情况也不可能发生。有人会问为什么不在这里直接调用MQ发送这个问题问得很好。如果你在同一个事务方法里既写数据库又发MQ本质上又回到了文章开头说的老路——发MQ失败的话你的数据库事务要不要回滚如果回滚MQ里可能已经收到消息了发送其实成功了但返回超时如果不回滚消息没发出去但业务数据一致只能靠后续补偿。所以事务内只管写表不碰任何远程调用这条铁律必须刻在项目规范里。2.3 事务边界控制别让Outbox表拖垮主业务使用事务注解时还有一个非常实际的调优点Outbox表要操作得足够快绝不能拖慢主事务。事务持续的时间越长数据库连接占用越久行锁和间隙锁的范围越大高并发下的吞吐量就下降得越快。所以Outbox表的插入动作必须轻量字段能少就少索引能精简就精简payload按JSON存储不要做多余的序列化处理。实践中我见过有人把payload写成TEXT存大段XML结果事务里做了大量字符串拼接接口RT直接翻倍。另外一个容易被忽略的地方是Outbox表的写入不要在事务里做复杂的查询或校验。业务数据检查应该在进入事务方法之前完成事务内只做必须的更新和插表。有人喜欢在事务里反复查库验证各种状态这会让锁持有时间变长在高并发下容易出现死锁和超时。3. 核心发送与兜底逻辑CronJob怎么把漏网之鱼捞回来Outbox表写进去了谁来把消息真正发到MQ这是第二个阶段的重点。业内常见的做法是使用一个独立的发送组件可以是常驻的轮询任务也可以是定时任务框架。我这里介绍两种方式一种是基于Spring自带的Scheduled做轻量CronJob另一种是引入分布式任务调度平台如XXL-Job做高可用补偿。3.1 常规发送线程基本的轮询发布器先用Spring Cache的Scheduled实现一个简单的发送Job负责把待发送的消息发出去。Component public class OutboxPublisher { private final OutboxEventMapper outboxEventMapper; private final RocketMQTemplate rockeTMqTemplate; private final ObjectMapper objectMapper; private static final int BATCH_SIZE 100; Scheduled(fixedDelay 2000) public void publishPendingEvents() { ListOutboxEvent pendingEvents outboxEventMapper.selectPendingEvents( OutboxStatus.PENDING.getCode(), BATCH_SIZE); for (OutboxEvent event : pendingEvents) { try { // 这里假设RocketMQ换成Kafka或RabbitMQ同理 SendResult result rockeTMqTemplate.syncSend( event.getBizType(), event.getPayload(), event.getEventId()); if (result.getSendStatus() SendStatus.SEND_OK) { outboxEventMapper.markSent(event.getId()); } } catch (Exception e) { // 记录失败次数等待下一轮或者CronJob兜底 outboxEventMapper.incrementRetry(event.getId()); log.error(发送Outbox事件失败, eventId{}, event.getEventId(), e); } } } }这段代码里有几个细节值得展开说。一个是fixedDelay 2000意思是上一次执行完毕后再等2秒才执行下一轮。为什么不用fixedRate因为fixedRate是固定频率不管上一次执行多久只要到了时间就触发下一次。如果批量发送遇到MQ抖动执行时间变长fixedRate会导致任务堆积线程池被打爆。fixedDelay可以天然规避这个问题。另一个是syncSend同步发送。有人为了追求性能会使用异步发送但在Outbox场景里我建议用同步发送。原因是消息可靠性才是这个方案的第一诉求性能损失可以通过限制批次大小、调整轮询间隔来弥补。一个批次发100条消息每条同步发送也就几毫秒甚至更短对整体影响完全可以接受。还有一个是失败处理。我在代码里只做了incrementRetry也就是把retry_count加一。这里要注意重试次数不能无限增长否则一个永远发不出去的坏消息会卡在表里反复占资源。设置一个阈值比如超过5次后把状态改成FAILED同时触发告警让值班人员来人工排查。3.2 兜底CronJob把滞留时间超过阈值的记录捞出来重发常规发送线程已经能覆盖99%的场景但还有1%的极端情况发送线程所在的服务重启、内存崩溃、JVM长时间Full GC导致一批消息在很长一段时间内无人问津。这时候就需要兜底CronJob出场了。兜底CronJob的思路是不管发送线程是否存在我都定期扫描那些“滞留太久”的消息。判断标准就是next_retry_time字段。Component public class OutboxCompensateJob { private final OutboxEventMapper outboxEventMapper; private final RocketMQTemplate rocketMQTemplate; Scheduled(cron 0 0/1 * * * ?) public void compensate() { // 捞出发送失败超过3次、且下次重试时间已经到了的记录 ListOutboxEvent staleEvents outboxEventMapper.selectStaleEvents( OutboxStatus.PENDING.getCode(), new Date(), 3, 200); for (OutboxEvent event : staleEvents) { try { SendResult result rocketMQTemplate.syncSend( event.getBizType(), event.getPayload(), event.getEventId()); if (result.getSendStatus() SendStatus.SEND_OK) { outboxEventMapper.markSent(event.getId()); } } catch (Exception e) { // 计算下次重试时间指数退避 Date nextRetry calculateNextRetryTime(event.getRetryCount()); outboxEventMapper.updateNextRetryTime(event.getId(), nextRetry); log.error(补偿任务发送失败, eventId{}, event.getEventId(), e); } } } }注意兜底CronJob的触发条件不能太激进。如果每一分钟就全表扫描一次对数据库压力很大而且可能和常规发送线程抢同一条记录造成重复发送。重复发送本身不可怕因为下游有幂等但无意义的重复扫描会浪费资源。我的建议是常规发送线程每2秒扫一次兜底CronJob每1分钟扫一次而且只扫next_retry_time小于当前时间、重试次数大于等于3的记录。这样两级任务各有分工常规发送线程负责把新鲜的事件尽快发出去兜底CronJob负责抢救那些在常规流程里卡住的异常数据。3.3 重试策略固定间隔和指数退避怎么选失败重试的策略直接影响系统的稳定性和数据的一致性。最简单的方案是固定间隔重试比如每5分钟重试一次最多重试10次。好处是逻辑简单坏处是如果下游MQ长时间不可用每次重试都会触发一轮无效的数据库扫描和远程调用浪费资源。更推荐的方案是指数退避加最大次数限制。第一次失败后等1分钟重试第二次失败等2分钟第三次等4分钟第四次等8分钟……直到达到上限比如64分钟或128分钟之后标记为人工处理。指数退避的好处是系统刚出问题时频繁重试能尽快追上新状态如果问题持续重试频率逐步降低相当于给系统一个“冷静期”。这个策略在分布式系统里非常常用RPC调用、消息消费失败重试都适用。private Date calculateNextRetryTime(int retryCount) { // 指数退避1分钟、2分钟、4分钟、8分钟……最多到64分钟 long backoffMillis Math.min( (long) Math.pow(2, retryCount) * 60 * 1000, 64L * 60 * 1000); return new Date(System.currentTimeMillis() backoffMillis); }这段代码的核心在Math.pow(2, retryCount)每次重试的间隔时间指数级增长最终封顶在64分钟避免无限拉长。4. 分布式部署场景下的关键问题与处理单体应用里用Outbox和CronJob已经非常成熟了但一旦把服务部署成多实例、上Kubernetes就会冒出一堆新的问题。这些问题是你在架构设计阶段就要提前想好的否则上线后一定会踩坑。4.1 多实例部署时定时任务重复执行怎么办假设你的订单服务部署了3个实例每个实例都会执行Scheduled任务。如果不做任何控制同一个Outbox事件可能被3个实例同时捞出来同时发到MQ下游同时收到3条一样的消息。即使下游有幂等也会造成大量无意义的带宽浪费和数据库压力。解决思路有几个。第一个思路是使用分布式锁。在发送Job执行前先尝试获取一个锁只有拿到锁的实例才能执行扫描和发送其他实例直接跳过本轮。可以用Redis的SETNX加过期时间实现也可以用Curator的分布式锁。public void publishWithLock() { String lockKey lock:outbox:publish; boolean locked redisTemplate.opsForValue() .setIfAbsent(lockKey, 1, Duration.ofSeconds(10)); if (!locked) { // 其他实例已经在执行本实例跳过 return; } try { publishPendingEvents(); } finally { redisTemplate.delete(lockKey); } }注意锁的过期时间要大于任务执行时间否则锁提前释放其他实例会趁虚而入。任务执行时间一般是几十毫秒到几百毫秒10秒的过期时间足够。第二个思路是使用数据库层面的唯一约束或行级锁。比如在Outbox表上额外维护一个lock_owner字段每个实例发送前先尝试UPDATE outbox_event SET lock_owner ? WHERE id ? AND lock_owner IS NULL影响行数为1说明抢锁成功为0说明被别人抢了。这种方式不需要额外引入Redis在数据库层面解决但要小心死锁和长事务。第三个思路更大胆——直接用分布式定时调度平台比如XXL-Job。它会自动保证一个任务在同一时刻只有一个实例执行不需要你操心锁的问题还自带调度日志、执行监控、失败告警等功能。如果你的公司已经有这类基础设施优先用它没必要重复造轮子。4.2 数据库和MQ之间的最终一致真的能保证不丢吗很多人问Outbox CronJob这套方案是不是能100%保证消息不丢严格来说它保证的是“业务数据入库的消息最终一定会被发送”。但这里有一个隐含前提CronJob本身不能挂。如果写入Outbox表的服务挂了并且重启后没有补扫机制那滞留的消息永远发不出去。所以真正的生产环境中还有一个隐藏角色对账监控。你要在系统里加一个监控任务统计outbox_event表里created_at超过某个阈值比如10分钟且status仍然不是已发送的记录数量。只要这个数量持续大于0就说明发送链路有异常需要告警到值班群。这个监控任务的设计思路和CronJob很像但职责完全不同。CronJob是“处理异常数据”监控任务是“发现异常并通知人”。两者配合起来才能形成一个相对完整的闭环事务写表保证本地原子定时任务保证消息发送监控告警保证问题及时暴露。4.3 消费端幂等Outbox不解决这个问题但你必须解决前面反复提到幂等这里展开说透。Outbox模式保证了消息“至少发送一次”但在MQ本身的机制下消费者收到的消息可能“晚到、重复、乱序”。这不是Outbox的锅而是分布式消息架构的常态。解决幂等最常见的手段是使用数据库唯一键去重。在消费端的业务处理表里建一个biz_id或event_id的唯一索引。消费者收到消息后先尝试插入一条“消费记录”如果插入成功说明这条消息是新的执行真正的业务逻辑如果插入冲突说明之前已经处理过直接跳过。Transactional public void handleOrderPaidEvent(OrderPaidMessage message) { try { // 尝试插入消费记录依靠唯一索引去重 consumeRecordMapper.insert(ConsumeRecord.builder() .eventId(message.getEventId()) .bizId(message.getOrderId()) .build()); } catch (DuplicateKeyException e) { // 已经消费过直接返回 log.info(重复消费跳过. eventId{}, message.getEventId()); return; } // 真正的业务逻辑加积分、发短信等 pointsService.addPoints(message.getUserId(), message.getPoints()); }这个方案的优点是不用依赖Redis等外部存储的原子性缺点是需要多建一张消费记录表。另外要注意插入消费记录和执行业务逻辑必须在同一个本地事务里否则可能出现消费记录插进去了但业务逻辑没执行的情况消息就被“假消费”了。5. 实战问题排查技巧Outbox方案的坑与解法关于这套方案网上讲原理的文章不少但真正上线后遇到的问题基本没人写。我基于自己的经验整理了几个高频问题供你参考。5.1 Outbox表积压严重、扫描性能下降表数据量过千万后SELECT ... WHERE status 0 AND next_retry_time NOW() ORDER BY id LIMIT 200这类查询会越跑越慢。即使有索引查询也还会扫描大量历史数据。解决方案有两个方向。一个是物理删除。已发送的记录定期清理比如只保留最近7天的数据更早的直接DELETE。配合定时任务在低峰期跑一次批量删除。另一个是归档。把已发送的记录转移到归档表或者冷存储业务热表只保留待发送和发送失败的记录。我的实际经验是Outbox表只做短期的状态存储10天前的数据如果不是为了对账基本没有保留价值。真要长期追溯把消息内容在MQ里多留几天就够了数据库里的记录可以清掉。5.2 消息重复既是问题也不是问题有的团队把Outbox方案上线后发现下游偶尔会收到两条一模一样的消息。查来查去原因可能是常规发送线程和兜底CronJob并发捞到了同一条记录也可能是发送方在调用MQ时超时了但实际已发送成功代码误判为失败继续重发。重复消息在这个方案里是预期内的事情所以处理重复的手段必须前置。除了前面说的消费记录表还有两种做法第一种是基于状态流转的幂等校验。下游维护一个业务状态的版本号只有当消息里的版本号大于当前版本号时才执行更新。第二种是使用状态机消息里携带目标状态下游判断当前状态能否流转到目标状态不能就丢弃。比较好的做法是同时使用消费记录表和业务状态校验双保险。只有消费记录表能挡住完全相同的重复挡不住不同消息到达顺序错乱导致的状态反转。5.3 事务内发MQ旧代码的顽疾很多团队引入Outbox时项目里已经积累了大量“事务内直接调MQ”的老代码。这些代码不改造还是会存在丢消息的隐患。改造有两种策略第一种存量消息先落库增量全部走Outbox。这个方法改动量小把老代码里同一个事务方法内发送MQ的逻辑摘出来改成先写表然后依赖后台任务统一发送。改完后回归测试有问题可以快速回滚。第二种逐步替换下游的消息主题。老主题继续保留新逻辑用新的主题名称下游进行双消费验证新链路没问题后直接把老代码下掉。这个策略适合消息体结构有大变化的情况。不管哪种策略改造过程中最忌讳的是“上线即删老代码”。正确的姿势是灰度观察至少一周看新链路的数据对比对无异常再清理旧逻辑。5.4 数据库连接池耗尽高并发下的边角问题Outbox模式多了一张表的插入操作在极高并发下会额外占用数据库连接。如果业务系统本身连接池配置偏小可能因为多出的写表操作导致连接不够用。解决方法是调大连接池的上限或者把Outbox表的插入操作放到专门的事务管理器里使用单独的数据源。后一种方案的实现会复杂一些但能有效隔离主链路和辅助链路。一般情况下调大连接池参数就能解决问题建议先用这个低成本方案。6. 上生产前的验收要点与个人体会文章最后把我自己上线这套方案时的验收清单分享一下可以看作一个检查表。第一个检查点写业务表和写Outbox表是否真的在同一个本地事务里。确认的方式很简单业务数据更新后主动抛一个异常看Outbox表里有没有落数据。如果Outbox表里也回滚了说明事务生效如果Outbox表里有残留说明事务边界没控制好。第二个检查点只停掉常规发送线程模拟CronJob兜底场景看滞留消息能否在下一个调度周期被重新发送。这个测试很重要能直接验证你的兜底链路是否真的可用。第三个检查点杀掉一个实例看另一个实例能否接续发送未完成的消息。如果你的分布式锁实现有问题杀实例后锁一直没释放消息就会无限期卡住。用这个操作来验证锁的过期时间是否合理。第四个检查点给下游消费逻辑做一次重复消息压测看幂等逻辑是否可靠。把同一条消息发送100遍看业务数据是否仍然只有一份。这几个检查都通过这套方案的质量基本就有保障了。从我个人的体会来说Outbox模式的核心价值不在于它有多先进而在于它把“分布式事务的一致性问题”巧妙地转换成了“本地事务 异步重试 幂等消费”三个简单问题的叠加。复杂问题被拆简单了系统的稳定性自然就上来了。如果你正在为消息丢失、分布式事务头疼别急着上那些分布式事务中间件先试试这个方案大多数时候它已经足够解决问题了。最后分享一个我踩过的坑Outbox表的event_id生成千万别用数据库自增ID因为消息落到MQ后下游拿到的ID如果和别的业务主键撞了幂等逻辑会非常混乱。老老实实用UUID或雪花ID这个字段承担的是全局唯一职责值得你多花几行代码。