1. 从一次线上事故说起为什么需要OutBox模式去年我们团队做了一套订单履约系统核心链路是用户下单后订单服务在本地数据库写入订单记录然后发一条消息到消息队列通知库存服务扣减库存、通知积分服务加积分、通知风控服务做审核。上线头两个月跑得挺稳直到某天大促运维同事发现有一批订单状态是已创建但库存没扣、积分没加用户投诉电话直接打爆了客服。排查下来原因很典型订单服务写完数据库之后在发消息那一步网络抖动导致消息发送超时代码里抛了异常本地事务回滚了——但问题是数据库写入和消息发送根本不在同一个事务里。更麻烦的是另一种情况数据库提交成功了消息也发出去了但消息队列的确认回执没收到代码重试又发了一遍结果库存被扣了两次。这就是分布式系统里最经典的本地事务与消息发送原子性问题。你没法让数据库和消息队列共享一个事务因为它们是两个独立的系统各自有各自的提交协议。硬要凑在一起要么引入两阶段提交这种重到不行的方案要么就得接受要么丢消息、要么重复消息的现实。OutBox模式就是在这个背景下被广泛采用的解法。它的核心思路特别朴素既然数据库和消息队列没法放在一个事务里那就把要发的消息先写进数据库跟业务数据在同一个本地事务里提交。然后另起一个异步的MessageRelay消息中继组件从数据库里把待发的消息捞出来投递到消息队列。这样业务数据和消息记录要么一起成功要么一起失败原子性由本地数据库事务保证消息的最终投递则由中继组件负责配合重试机制实现最终一致性。这篇文章我会把OutBox模式从设计思路到落地实现完整拆一遍包括表结构怎么设计、MessageRelay怎么写、幂等怎么做、常见坑怎么避。适合正在做微服务、事件驱动架构或者被消息发不出去/发重复了折磨过的后端同学。读完你至少能拿到一套可以直接抄作业的方案。2. OutBox模式整体设计与核心思路拆解2.1 核心问题拆解原子性到底难在哪先把问题说透。假设订单服务要做的两件事是在order表插入一条订单记录往消息队列发一条OrderCreated事件。理想情况下这两件事要么都成功要么都失败。但现实是它们分属两个系统中间隔着网络。我们逐一分析几种失败场景场景数据库写入消息发送结果正常成功成功一致写库后发消息前宕机成功未发送消息丢失下游不知道有新订单发消息后写库失败未提交已发送下游收到不存在的订单脏数据发消息成功但回执丢失成功实际已发送重试导致重复消息传统做法是先写库再发消息发失败就重试但这解决不了写库成功、进程崩溃、消息没发出去的情况。也有人用本地消息表定时补偿其实那就是OutBox模式的雏形。2.2 OutBox的核心思想把消息当成业务数据的一部分OutBox模式的精妙之处在于视角转换不要把发消息看成一个独立的网络动作而是把它看成一次数据库写入。具体来说在业务库里建一张outbox表或者叫message_outbox、event_outbox字段大致包括消息ID、聚合根ID、事件类型、消息体、状态、创建时间、重试次数等。业务代码在同一个本地事务里做两件事BEGIN; INSERT INTO orders (id, user_id, amount, status) VALUES (...); INSERT INTO outbox (id, aggregate_id, event_type, payload, status, created_at) VALUES (..., OrderCreated, {orderId:...}, PENDING, NOW()); COMMIT;这两条INSERT在同一个数据库事务里数据库的ACID特性保证了它们要么一起成功要么一起回滚。消息此时并没有真正发到消息队列只是待发送状态躺在数据库里。然后有一个独立的MessageRelay进程可以是独立服务也可以是应用内的后台线程不断轮询outbox表里状态为PENDING的记录逐条投递到消息队列。投递成功后把状态更新为SENT或直接删除记录。投递失败就保留PENDING下一轮继续重试。2.3 为什么这个方案能成立关键点在于消息的产生和消息的投递被解耦了。消息的产生跟业务数据在同一个本地事务里原子性由数据库保证绝对不会出现业务成功但消息没记录的情况。消息的投递由MessageRelay异步完成允许失败重试只要最终投递成功即可这就是最终一致性。有人会问那MessageRelay挂了怎么办消息不就一直躺在数据库里答案是MessageRelay是无状态的可以多实例部署挂了重启继续捞。只要数据库里的PENDING记录还在消息就不会丢。这比内存里存着待发消息、进程一挂全没了可靠得多。2.4 方案选型对比OutBox vs 其他方案在决定用OutBox之前我对比过几种常见方案这里列个表方便你判断方案原子性保证复杂度消息延迟适用场景先写库再发消息无低低能容忍少量丢失的场景本地消息表OutBox强中秒级大多数业务系统事务消息RocketMQ强中低已用RocketMQ且能接受其事务API两阶段提交XA强高高极少用性能差最大努力通知弱低不定对账类场景OutBox的优势是不依赖特定消息中间件MySQL、PostgreSQL都能用实现成本可控可靠性有保障。缺点是消息投递有秒级延迟取决于轮询间隔以及对数据库有一点点额外写入压力。对于绝大多数业务系统这个代价完全可接受。提示如果你的系统已经在用RocketMQ它的事务消息机制本质上也是一种OutBox思想半消息回查可以直接用。但如果消息中间件不支持事务消息OutBox就是最通用的选择。3. 核心细节解析与实操要点3.1 OutBox表结构设计的关键字段表结构设计直接决定了后续MessageRelay好不好写、性能扛不扛得住。我踩过几次坑之后总结出一套比较稳的字段设计CREATE TABLE outbox ( id BIGINT NOT NULL AUTO_INCREMENT, message_id VARCHAR(64) NOT NULL COMMENT 全局唯一消息ID用于幂等, aggregate_type VARCHAR(64) NOT NULL COMMENT 聚合类型如Order, aggregate_id VARCHAR(64) NOT NULL COMMENT 聚合根ID如orderId, event_type VARCHAR(64) NOT NULL COMMENT 事件类型如OrderCreated, payload TEXT NOT NULL COMMENT 消息体JSON格式, status TINYINT NOT NULL DEFAULT 0 COMMENT 0PENDING,1SENT,2FAILED, retry_count INT NOT NULL DEFAULT 0 COMMENT 重试次数, next_retry_at DATETIME NOT NULL DEFAULT NOW() COMMENT 下次重试时间, created_at DATETIME NOT NULL DEFAULT NOW(), updated_at DATETIME NOT NULL DEFAULT NOW() ON UPDATE NOW(), PRIMARY KEY (id), UNIQUE KEY uk_message_id (message_id), KEY idx_status_next_retry (status, next_retry_at) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;几个字段的设计意图值得展开说message_id全局唯一建议用UUID或者业务ID事件类型时间戳组合。它的作用是让下游消费者做幂等——同一条消息即使被投递多次消费者根据message_id去重只处理一次。这是OutBox模式必须配套的机制因为MessageRelay的重试天然会导致重复投递。status next_retry_at 联合索引这是MessageRelay轮询的核心索引。查询语句是WHERE status0 AND next_retry_at NOW() ORDER BY id LIMIT 100走这个联合索引效率很高。为什么要next_retry_at而不是直接按created_at因为失败重试需要退避不能让失败的消息每秒钟都被捞出来重试那样会把数据库打爆。payload用TEXT而不是JSON类型MySQL 5.7之前没有原生JSON类型用TEXT兼容性更好。即使有JSON类型我也倾向于TEXT因为消息体格式可能变化TEXT更灵活序列化反序列化在应用层做。3.2 MessageRelay的三种实现形态MessageRelay是OutBox模式的心脏它的实现方式直接决定了系统的可靠性和复杂度。我实践过三种形态各有适用场景形态一应用内定时任务。在业务服务里起一个Scheduled定时任务每隔几百毫秒扫一次outbox表。优点是部署简单不用额外组件缺点是跟业务服务耦合业务服务多实例部署时会有多个Relay同时跑需要加分布式锁或者用SELECT ... FOR UPDATE SKIP LOCKED来避免重复捞取。形态二独立中继服务。把MessageRelay拆成一个独立的微服务专门负责扫表投递。优点是职责清晰可以独立扩缩容业务服务不用管消息投递缺点是多了一个服务要运维。形态三CDC变更数据捕获。用Debezium之类的工具监听数据库binlog捕获outbox表的INSERT事件直接投递到消息队列。优点是实时性极高几乎无延迟缺点是引入了CDC组件运维复杂度上升而且binlog解析对表结构变化比较敏感。我个人的选择是中小规模用形态一配合SKIP LOCKED规模大了或者团队有运维能力上形态二对延迟极其敏感的场景才考虑形态三。下面重点讲形态一和形态二的实现细节。3.3 并发捞取SKIP LOCKED是关键多实例部署时最怕的是两个Relay实例捞到同一条消息重复投递。传统做法是用SELECT ... FOR UPDATE加行锁但这样会导致实例之间互相阻塞吞吐上不去。MySQL 8.0和PostgreSQL 9.5之后支持SKIP LOCKED完美解决这个问题SELECT * FROM outbox WHERE status 0 AND next_retry_at NOW() ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED;SKIP LOCKED的含义是如果某行已经被其他事务锁住就跳过它不等待。这样多个Relay实例可以并行捞取不同的消息互不阻塞。实测下来4个实例并行处理吞吐能提升3倍以上。注意MySQL 8.0以下不支持SKIP LOCKED如果你还在用5.7要么升级要么用乐观锁状态机的方式先UPDATE outbox SET status1 WHERE id? AND status0根据affected rows判断是否抢到抢到再投递。这种方式也能避免重复但多一次UPDATE。3.4 幂等设计消费端必须做的事OutBox模式保证了消息不丢但不保证消息不重复。MessageRelay投递成功但更新状态失败、或者投递超时后重试都会导致同一条消息被投递多次。所以消费端必须做幂等。幂等的常见做法有三种唯一索引去重消费端建一张consumed_message表message_id做唯一索引处理前先INSERT冲突就跳过。简单可靠适合大多数场景。业务状态机比如订单状态从待支付到已支付只有当前状态是待支付时才处理天然幂等。Redis去重用message_id做keySETNX成功才处理设置合理过期时间。性能好但Redis挂了可能丢去重记录。我一般推荐第一种因为它是持久化的最可靠。第二种需要业务逻辑配合不是所有场景都适用。第三种适合对性能要求极高、能容忍极端情况下少量重复的场景。4. 实操过程与核心环节实现4.1 业务侧写入一个事务搞定两件事先看业务代码怎么写。以Spring Boot MyBatis为例核心是把订单插入和outbox插入放在同一个Transactional方法里Service public class OrderService { Autowired private OrderMapper orderMapper; Autowired private OutboxMapper outboxMapper; Transactional(rollbackFor Exception.class) public void createOrder(CreateOrderCmd cmd) { // 1. 写业务数据 Order order new Order(); order.setId(cmd.getOrderId()); order.setUserId(cmd.getUserId()); order.setAmount(cmd.getAmount()); order.setStatus(CREATED); orderMapper.insert(order); // 2. 写outbox消息同一个事务 OutboxMessage msg new OutboxMessage(); msg.setMessageId(UUID.randomUUID().toString()); msg.setAggregateType(Order); msg.setAggregateId(cmd.getOrderId()); msg.setEventType(OrderCreated); msg.setPayload(JSON.toJSONString(buildEvent(order))); msg.setStatus(0); msg.setNextRetryAt(new Date()); outboxMapper.insert(msg); } }这里有个细节message_id的生成。我建议在业务侧生成而不是让数据库自增。因为消息ID需要在投递前就确定下游消费端才能用它做幂等。用UUID最简单如果嫌UUID太长影响索引性能可以用雪花算法生成Long型ID再转字符串。另一个细节payload里放什么。我的经验是放完整的领域事件包括事件发生的时间、聚合根ID、变更的关键字段。不要只放一个ID让消费端回查那样消费端还得调你的接口耦合太深。事件驱动架构的精髓就是事件自包含。4.2 MessageRelay的完整实现下面是一个基于Spring定时任务SKIP LOCKED的MessageRelay实现可以直接参考Component public class MessageRelay { Autowired private OutboxMapper outboxMapper; Autowired private KafkaTemplateString, String kafkaTemplate; private static final int BATCH_SIZE 100; private static final int MAX_RETRY 5; Scheduled(fixedDelay 500) public void relay() { ListOutboxMessage messages outboxMapper.selectPendingForUpdate(BATCH_SIZE); for (OutboxMessage msg : messages) { try { // 投递到消息队列key用aggregateId保证同聚合有序 kafkaTemplate.send(order-events, msg.getAggregateId(), msg.getPayload()) .get(3, TimeUnit.SECONDS); // 投递成功标记为已发送 outboxMapper.markSent(msg.getId()); } catch (Exception e) { // 投递失败增加重试次数计算下次重试时间 int retry msg.getRetryCount() 1; if (retry MAX_RETRY) { outboxMapper.markFailed(msg.getId()); log.error(消息投递失败超过最大重试次数, messageId{}, msg.getMessageId(), e); } else { // 指数退避2^retry 秒 long delaySeconds (long) Math.pow(2, retry); outboxMapper.scheduleRetry(msg.getId(), retry, delaySeconds); } } } } }对应的Mapper SQLselect idselectPendingForUpdate resultTypeOutboxMessage SELECT * FROM outbox WHERE status 0 AND next_retry_at lt; NOW() ORDER BY id LIMIT #{batchSize} FOR UPDATE SKIP LOCKED /select update idmarkSent UPDATE outbox SET status 1, updated_at NOW() WHERE id #{id} /update update idscheduleRetry UPDATE outbox SET retry_count #{retry}, next_retry_at DATE_ADD(NOW(), INTERVAL #{delaySeconds} SECOND), updated_at NOW() WHERE id #{id} /update4.3 参数计算轮询间隔和批量大小怎么定这两个参数直接决定消息延迟和数据库压力需要算一下。轮询间隔假设业务峰值QPS是1000每条业务操作产生1条消息那么每秒产生1000条待发消息。如果轮询间隔是500ms每轮要处理500条。批量大小设100的话一轮处理不完会积压。所以要么缩小间隔要么增大批量。我的经验公式是批量大小 峰值QPS × 轮询间隔(秒) / Relay实例数。按上面的例子1000 × 0.5 / 4 125所以批量设150比较稳妥。延迟消息从产生到投递的平均延迟约等于轮询间隔 / 2。500ms间隔对应平均250ms延迟对大多数业务够用。如果要求100ms以内间隔要设到200ms以下但数据库压力会上升。数据库压力每轮一次SELECT N次UPDATE。1000 QPS下每秒约2次SELECT 1000次UPDATE。这个量级对MySQL来说不算大但要注意outbox表要定期清理已发送的记录否则表会越来越大查询变慢。4.4 清理策略已发送消息怎么处理status1的记录不能一直留着否则表膨胀。三种清理方式定时删除每天凌晨删除7天前的status1记录。简单但删除大批量数据可能锁表建议分批删。分区表按天分区直接DROP老分区。性能最好但需要DBA配合。归档到历史表先INSERT到outbox_history再DELETE。适合需要保留审计记录的场景。我一般用第一种分批删除每批1000条加个LIMITDELETE FROM outbox WHERE status 1 AND created_at DATE_SUB(NOW(), INTERVAL 7 DAY) LIMIT 1000;循环执行直到affected rows为0。这样不会长时间锁表。5. 常见问题与排查技巧实录5.1 消息积压了怎么办现象outbox表里status0的记录越来越多消息延迟飙升。排查思路分三步第一步看Relay是否在跑。检查Relay服务的日志和线程状态确认定时任务没被阻塞。我遇到过一次是Relay里调用的下游接口超时导致整个批次卡住后面的消息全积压。解决办法是给投递加超时单条消息投递失败不影响其他消息。第二步看数据库是否有锁等待。用SHOW PROCESSLIST或者information_schema.innodb_trx看有没有长事务。如果业务侧有慢事务一直不提交Relay的FOR UPDATE会等待导致积压。第三步看消息队列是否健康。如果Kafka集群有问题投递全部失败消息会进入重试队列next_retry_at被推后表现为积压。这时候要优先恢复消息队列。5.2 消息重复投递怎么定位现象下游消费端报重复处理或者数据被改了两次。先确认是不是OutBox导致的重复。查outbox表看同一条message_id是否有多次投递记录。如果status从0变1又变回0理论上不会但如果有并发bug可能或者Relay日志里同一条消息被投递多次那就是Relay的问题。更常见的情况是Relay投递成功但markSent更新失败比如数据库连接闪断下一轮又捞出来投递一次。这种重复是OutBox模式的固有特性必须靠消费端幂等兜底。所以我在3.4节强调消费端幂等是必须的不是可选的。排查时可以在消费端加日志记录每次处理的message_id看是否有重复。如果有检查消费端的幂等逻辑是否生效。5.3 常见问题速查表问题现象可能原因排查方法解决方案消息延迟高轮询间隔大/Relay实例少看outbox积压量缩小间隔、增加实例、增大批量消息丢失业务事务未提交/Relay未启动查outbox是否有PENDING记录检查事务注解、Relay健康状态消息重复Relay重试/消费端无幂等查消费端日志message_id消费端加唯一索引去重outbox表膨胀未清理已发送记录查表行数定时分批删除Relay CPU高轮询过频/批量过大看Relay监控调整间隔和批量投递顺序乱多实例并发投递查消息队列分区策略用aggregateId做分区key5.4 几个我踩过的坑坑一message_id用数据库自增。一开始图省事让outbox表的id自增做主键message_id也用它。结果业务侧在插入前拿不到ID没法在payload里带上消费端做幂等时只能等插入后回查。后来改成业务侧生成UUID问题解决。坑二payload里放了大对象。有次把整个订单对象序列化进去包括几十个字段单条消息几十KBoutbox表迅速膨胀查询变慢。后来精简成只放消费端需要的字段控制在1KB以内。坑三Relay和业务共用一个数据库连接池。Relay的轮询和投递会占用连接高峰期把连接池占满业务请求拿不到连接。后来给Relay单独配了一个小连接池隔离开了。坑四忘记处理投递超时。Kafka的send是异步的如果不调.get(timeout)消息可能还在缓冲区就返回了Relay以为投递成功实际没发出去。一定要同步等待确认或者用带回调的异步方式并正确处理回调。提示如果你的消息队列支持批量发送可以把一批消息打包发送减少网络往返。但要注意批量发送的失败处理——部分成功部分失败时要能精确知道哪些失败了只重试失败的。6. 进阶优化与扩展思路6.1 用CDC替代轮询降低延迟轮询方案的天花板是轮询间隔想做到毫秒级延迟就得用CDC。Debezium监听MySQL binlog捕获outbox表的INSERT事件直接投递到Kafka。这样消息产生后几乎立刻投递延迟从几百毫秒降到几十毫秒。代价是运维复杂度上升要部署Debezium、Kafka Connect要处理binlog格式变化、连接器重启等问题。我的建议是延迟要求不高的场景别上CDC轮询够用且简单可靠。真到了需要CDC的规模团队一般也有相应的运维能力了。6.2 多租户场景下的隔离如果系统是多租户的outbox表里混了所有租户的消息Relay投递时要注意隔离。一种做法是加tenant_id字段Relay按租户分组投递避免一个租户的消息积压影响其他租户。另一种是每个租户独立的outbox表隔离更彻底但管理成本高。6.3 与Saga模式的配合OutBox模式经常和Saga模式一起用。Saga的每个步骤都是一个本地事务每个步骤完成后发一个事件通知下一步。这个发事件的动作就用OutBox来保证原子性。这样整个Saga的每一步都是可靠的配合补偿事务实现分布式事务的最终一致性。举个例子订单Saga分三步——创建订单、扣库存、扣款。每步完成后通过OutBox发事件触发下一步。如果扣款失败发一个补偿事件触发库存回滚和订单取消。整个链路的消息可靠性都由OutBox保证。6.4 监控指标该看什么上线OutBox之后这几个指标必须监控outbox积压量SELECT COUNT(*) FROM outbox WHERE status0超过阈值告警。消息投递延迟消息创建时间到投递成功时间的差值P99超过1秒告警。投递失败率失败次数/总投递次数超过1%告警。Relay处理速率每秒处理消息数突然下降要排查。outbox表大小行数和磁盘占用增长过快说明清理没生效。这些指标用Prometheus Grafana就能搭起来Relay里埋点上报即可。我在实际项目里把这套OutBox方案跑了两年多经历过双十一级别的流量消息零丢失重复率控制在万分之一以下靠消费端幂等兜底。最大的体会是分布式系统里没有银弹OutBox也不是万能的它用消息延迟换可靠性用消费端幂等换投递重复。你得清楚自己系统能接受什么样的延迟、能不能做幂等再决定用不用。如果业务对延迟极其敏感又做不了幂等那OutBox可能不适合你得考虑别的方案。但对绝大多数业务系统来说这套方案的性价比是最高的实现成本可控可靠性有保障值得作为标准基础设施沉淀下来。
阅读完成 · 觉得有帮助?