事务发件箱模式:可靠消息投递的工程解法

深入事务发件箱模式,解决业务库与消息队列双写一致性问题,对比轮询发布与事务日志监听两种方式,结合 CDC 联动实现幂等可靠的最终一致消息投递

在微服务与事件驱动架构中,最隐蔽的故障往往不是业务逻辑错误,而是「数据库已提交、消息却丢失」的幽灵。事务发件箱(Transactional Outbox)用一个本地事务加一张业务表,就化解了双写不一致这个经典难题。本文从模式原理出发,对比轮询发布与事务日志监听两种实现,并结合 CDC 与幂等消费给出生产级落地细节。

前置基础可先阅读 Spring 集成消息队列 与 分布式事务:Seata、TCC 与 Saga 模式详解。

1. 双写一致性困境

1.1 一个订单引发的数据分裂

业务系统最常见的需求是「落库之后发事件」。订单表写入与消息发送是两个独立系统,任何一次失败都会造成二者不一致:

失败场景结果
先发消息后写库,写库失败下游收到事件但库里没有订单
先写库后发消息,发送失败订单已存在但下游无感知
发送超时但实际已送达重发造成重复消息
库已回滚但消息已发送幽灵事件
// 典型但错误的双写
@Transactional
public Order createOrder(OrderRequest req) {
    Order order = orderDao.save(req.toOrder());      // ① 写库
    eventPublisher.publish(new OrderCreated(order));  // ② 发消息
    // 若第 ② 步失败,① 已提交,事务无法回滚消息
    return order;
}

1.2 为什么不用本地事务保护消息发送

本地事务只能保护同一资源(同一个数据库连接),而消息队列是外部资源。把 kafkaTemplate.send 放进 @Transactional 事务里,依赖的是消息中间件参与 JTA 分布式事务,成本和复杂度都极高,且大多数场景下根本不值得。

1.3 发件箱模式的核心思想

把「发消息」这个外部副作用,替换为「往本地表插一行」这个本地副作用。二者共享同一个数据库事务,要么都成功,要么都回滚:

业务方法(一个本地事务):
  ① 更新业务表(订单状态、余额扣减)
  ② 往 outbox 表插入事件记录
  → 提交
后台发布器:
  ③ 扫描 outbox 表中未发布的事件
  ④ 发送到消息队列
  ⑤ 标记为已发布

2. 发件箱表设计

2.1 表结构与 DDL

发件箱表是模式的载体,字段设计要兼顾发布器查询效率与消息体存储:

CREATE TABLE outbox_event (
    id            BIGINT AUTO_INCREMENT PRIMARY KEY,
    aggregate_id  VARCHAR(64)  NOT NULL COMMENT '聚合根 ID,如订单号',
    aggregate_type VARCHAR(32) NOT NULL COMMENT '聚合类型,如 ORDER',
    event_type    VARCHAR(64)  NOT NULL COMMENT '事件类型,如 ORDER_CREATED',
    payload       JSON         NOT NULL COMMENT '事件体,序列化后的业务数据',
    status        TINYINT      NOT NULL DEFAULT 0 COMMENT '0待发布 1已发布 2失败',
    retry_count   INT          NOT NULL DEFAULT 0,
    created_at    DATETIME(3)  NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
    published_at  DATETIME(3)  DEFAULT NULL,
    KEY idx_ready (status, id)
) ENGINE=InnoDB;

2.2 事件写入与本地事务

业务方法与发件箱写入必须在同一个事务内,顺序通常是先更新业务表、后插入事件:

@Service
@Slf4j
public class OrderApplicationService {

    private final OrderDao orderDao;
    private final OutboxDao outboxDao;

    @Transactional
    public Order createOrder(OrderRequest req) {
        Order order = orderDao.save(req.toOrder());
        outboxDao.insert(OutboxEvent.builder()
                .aggregateId(order.getId())
                .aggregateType("ORDER")
                .eventType("ORDER_CREATED")
                .payload(JsonUtils.toJson(order))
                .build());
        return order;
    }
}

2.3 索引与写入优化

idx_ready (status, id) 让发布器能用最少的回表找到待发布事件。写入侧要保持轻量:payload 存序列化后的最终状态即可,不要在其中嵌入大对象。对极高吞吐场景,可对 outbox_event 做按事件类型的分区或分表。

3. 发布方式一 轮询发布器

3.1 分批查询与状态机

轮询发布器是发件箱模式最朴素可靠的实现:定时扫描未发布记录,发送成功后更新状态。状态机要覆盖「成功、失败、重试、死信」四种情况。

3.2 发布器实现

@Component
@Slf4j
public class OutboxPoller {

    private final OutboxDao outboxDao;
    private final KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 1000)
    public void publishReadyEvents() {
        List<OutboxEvent> batch = outboxDao.findReady(100);      // status=0 按 id 升序
        for (OutboxEvent event : batch) {
            publishOne(event);
        }
    }

    private void publishOne(OutboxEvent event) {
        try {
            kafkaTemplate.send(event.getEventType(), event.getPayload())
                    .get(5, TimeUnit.SECONDS);
            outboxDao.markPublished(event.getId());               // 幂等标记
        } catch (Exception e) {
            outboxDao.markFailed(event.getId());                  // 计失败
            log.error("outbox publish failed, id={}", event.getId(), e);
        }
    }
}

3.3 延迟、批量与有序性

  • 延迟:fixedDelay 与轮询批次共同决定延迟上限,秒级即可满足大多数事件驱动场景。
  • 吞吐:批次大小可调,findReady 加 LIMIT 防止单次处理过多。
  • 有序性:对同一 aggregate_id 的事件,发布器必须按 id 升序处理,且保证同一聚合的事件不被并发发布器并发发送。
// 保证有序:查询按 (aggregate_id, id) 排序,发布器单实例
// 或用分区锁:同一 aggregate_id 路由到同一发布线程

4. 发布方式二 事务日志监听器

4.1 从日志感知变更

轮询需要额外的定时任务与查询,而事务日志监听器直接消费数据库的变更日志(如 MySQL binlog、PostgreSQL WAL),零侵入地发现「新插入的 outbox 行」,即 CDC 思路。

// 事务监听器方案(借助 Debezium 或 Canal 读取 binlog)
// 当 outbox_event 出现 INSERT 时,实时解析行数据并投递到消息队列
// 优点:毫秒级延迟、无轮询查询开销
// 缺点:需要维护 CDC 组件,架构复杂度上升

4.2 事务边界内的本地监听器

若发布事件本身不强依赖消息队列,也可以借助 Spring 的 TransactionSynchronization,在事务提交后再发送,避免事务未提交就发消息导致下游读到旧数据:

@Transactional
public void createOrderWithEvent(OrderRequest req) {
    orderDao.save(req.toOrder());
    TransactionSynchronizationManager.registerSynchronization(
        new TransactionSynchronization() {
            @Override
            public void afterCommit() {
                eventPublisher.publish(new OrderCreated(req));   // 提交后才发送
            }
        });
}

注意:afterCommit 只解决「提交后再发」的时序,并不解决「发送失败」的可靠性,生产环境仍应配合发件箱表兜底。

4.3 两种发布方式对比

维度轮询发布器事务日志监听器
延迟秒级(取决于轮询间隔)毫秒级
复杂度低,仅需定时任务高,需 CDC 组件
对业务表侵入无无
数据库压力定期扫描索引binlog 消费,几乎无压力
适合场景中小流量、运维简单高吞吐、低延迟要求

5. 与 CDC 联动

5.1 Debezium 发件箱事件路由器

Debezium 官方提供 outbox-event-router 单消息转换(SMT),专门从 outbox 表读取变更并重组为规范的事件消息,是「事务日志监听 + 发件箱」的黄金组合。

transforms: outbox
transforms.outbox.type: io.debezium.transforms.outbox.EventRouter
transforms.outbox.table.fields.additional.placement: aggregate_id:header:aggregateId
transforms.outbox.table.fields.additional.placement: event_type:header:eventType
transforms.outbox.route.by.field: event_type
transforms.outbox.route.topic.regex: (?<routed>.+)

5.2 从发件箱表到 Kafka 的链路

MySQL binlog → Debezium Connector → outbox SMT 重塑
  → Kafka Topic(按 event_type 路由,带聚合 ID 头)
  → 下游消费者按 header 去重

这组链路把「谁去读发件箱表」完全委托给数据库日志,业务代码不再关心发布细节。

5.3 水位线与健康度

无论哪种发布方式,都要监控发件箱积压。核心指标是「未发布事件的最早时间」与「每分钟新增量」,它们直接反映事件投递的健康程度。

6. 幂等消费与重复消息处理

6.1 消费端幂等是最后防线

发件箱保证「至少一次」投递,重复不可避免:发布后标记失败会重发、下游处理超时会重投。消费端必须自身幂等。

@Component
@Slf4j
public class OrderEventConsumer {

    private final IdempotencyDao idemDao;
    private final OrderProjectionService projection;

    @KafkaListener(topics = "ORDER_CREATED", groupId = "order-projector")
    public void onOrderCreated(ConsumerRecord<String, String> record) {
        String eventId = record.headers().lastHeader("eventId").value() != null
                ? new String(record.headers().lastHeader("eventId").value())
                : record.key();
        if (!idemDao.tryAcquire("ORDER_CREATED", eventId)) {
            return;                                            // 已处理过,直接跳过
        }
        projection.buildProjection(record.value());            // 业务处理
        idemDao.commit("ORDER_CREATED", eventId);              // 标记完成
    }
}

6.2 去重表与唯一键

CREATE TABLE event_idempotency (
    event_key  VARCHAR(128) PRIMARY KEY,      -- 聚合类型 + 事件 ID
    process_at DATETIME(3) NOT NULL,
    status     TINYINT NOT NULL DEFAULT 1     -- 1处理中 2完成
) ENGINE=InnoDB;

tryAcquire 用 INSERT ... ON DUPLICATE KEY UPDATE 原子占位,多个并发消费者只有一个能抢到处理权。

6.3 乱序与重试策略

  • 乱序:同一聚合的事件必须保证按序消费,最简做法是发送时按 aggregate_id 做 key,Kafka 同 key 进同分区即天然有序。
  • 重试:消费失败不回滚占位、指数退避重试,超过阈值转人工。
  • 业务幂等:消息幂等之外,业务本身也应可重入,例如状态机校验「已支付订单拒绝重复扣款」。

7. 生产实践

7.1 表膨胀与归档

发件箱表只保留「待发布」数据才有价值,已发布行应定期归档:

-- 归档策略:每小时删除已发布超过 24 小时的数据
DELETE FROM outbox_event
 WHERE status = 1 AND published_at < NOW() - INTERVAL 1 DAY
 LIMIT 5000;

对超大数据量,可直接把已发布行迁入归档表或按时间分区后 DROP PARTITION。

7.2 监控与告警

指标告警阈值说明
未发布事件最早时间超过 5 分钟发布器挂掉或队列积压
未发布事件数量超过 1000大量事件堆积
发布失败率超过 5%消息队列异常
消费处理延迟超过 1 分钟下游消费缓慢
// 指标埋点示例
Gauge.builder("outbox.pending.max.age", pendingAgeSupplier)
        .register(registry);
Counter.builder("outbox.published.total").register(registry);

7.3 与分布式事务框架的关系

发件箱本质是最终一致性方案,适用「容忍短暂延迟」的异步链路;而 Seata AT/TCC 解决的是需要强一致的同步场景。两者选型边界清晰:订单主链路需要强一致时选 Seata,旁路事件通知优先发件箱。此外发件箱表还可作为对账依据,为下游重放提供数据源。

8. 总结

主题核心要点
双写困境业务库与消息队列无法被同一本地事务保护
模式本质把发消息替换为写本地发件箱表,共享事务
发布方式轮询发布器低延迟保障,事务日志监听毫秒级
CDC 联动Debezium outbox SMT 把发布职责交给 binlog
幂等消费去重表 + 唯一键,至少一次投递下保证精确处理
运维治理归档已发布数据,监控积压与失败率

事务发件箱用最小的成本把「业务提交」与「消息投递」绑定在同一个本地事务里,换来的是可靠的最终一致。它不要求消息中间件参与分布式事务,也不引入复杂的协调器,配合 CDC 与幂等消费,就能构建一条吞吐高、延迟可控、可观测的事件发布链路。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. JPMS 模块化:module-info 与 JLink 精简运行时
  2. CDC 数据同步:Debezium 与 Kafka 架构实战
  3. 可观测性工程:Micrometer 指标模型与 OTLP 导出