Skip to content

微服务数据一致性:最终一致性 + 本地消息表 vs 事务消息 vs 事件溯源

问题

微服务架构中,订单服务创建订单后必须通知库存服务扣减库存。如果服务间网络超时或库存服务宕机,订单数据写入了但库存没扣,如何保证最终一致?本地消息表和事务消息的核心区别在哪?事件溯源又是什么,什么场景才值得用?

直接抛结论:分布式事务(2PC/XA)理论上能保证强一致性,但实际生产中没人用——协调者单点、锁资源时间太长、性能差到没法接受。业界选择的是最终一致性:接受短暂的不一致状态,通过补偿机制确保最终数据是对的。

让我先画一条时序图,把三种方案的核心流程对齐,再逐个拆。

// 三种方案的核心流程对比

// 本地消息表
OrderService DB          MessageJob    Kafka    StockService
    |         |              |           |           |
    |--(1) insert order ----|           |           |
    |--(2) insert msg(status=0)         |           |
    |         |              |           |           |
    |         |   --(3) poll msg(status=0)          |
    |         |   --(4) send to kafka --|           |
    |         |              |           |--(5) consume
    |         |              |           |           |--(6) deduct stock
    |         |              |   --(7) ack/callback |
    |         |   --(8) update status=2  |           |

// 事务消息 (RocketMQ)
OrderService          RocketMQ Broker    StockService
    |                       |                 |
    |--(1) send half msg --|                 |
    |--(2) local tx: insert order            |
    |--(3) commit/rollback                   |
    |             |--(4) check callback       |  (如果2宕机了)
    |--(5) reply  |                          |
    |             |--(6) deliver msg --------|
    |             |                          |--(7) deduct stock

// 事件溯源
OrderService(EventStore)                  StockService
    |                                         |
    |--(1) append OrderCreated event          |
    |--(2) event bus publish -----------------|
    |                                         |--(3) replay & deduct
    |                                         |--(4) append StockDeducted event

方案一:本地消息表(Local Message Table)

核心思路:把"发消息"和"写业务数据"放在同一个本地事务里。

订单服务创建订单时,在自己的数据库里同时插入两条记录:一条订单表,一条消息表(状态为"待发送")。这两个操作在同一个本地事务中,要么都成功要么都回滚。后台定时任务轮询消息表,把未发送的消息推到 MQ,消费者消费成功后回调确认。

关键链路:

  1. 本地事务保证了业务数据和消息的原子性
  2. 定时任务保证了消息"至少一次"的投递
  3. 消费端幂等表保证了"恰好一次"的处理

实际踩坑:某电商双十一消息表写爆

我参与过一个电商项目,双十一当天本地消息表每秒写入 8000+ 条消息,定时任务每 5 秒扫一次,一次查 100 条。结果:

  • 消息表膨胀到 500 万行,idx_status_next_retry 索引失效,扫一次 3 秒
  • 订单都处理完了,库存还没扣,用户付了钱但发不出货
  • 尾段时间把扫描间隔缩到 1 秒,一次查 500 条,数据库连接池直接打满

教训:本地消息表方案必须做分表 + 单独的扫描线程池,不能和业务 API 共享连接池。

代码实现

1. 消息表设计(带分表键)

sql
-- 按 order_id 取模分 16 张表,每张表独立
CREATE TABLE local_message_%d (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  msg_id VARCHAR(64) NOT NULL UNIQUE COMMENT '全局唯一消息ID',
  biz_key VARCHAR(64) NOT NULL COMMENT '业务唯一键,用于幂等',
  content TEXT NOT NULL COMMENT '消息体JSON',
  status TINYINT NOT NULL DEFAULT 0 COMMENT '0=待发送 1=发送中 2=已发送 3=已到达最大重试',
  retry_count INT NOT NULL DEFAULT 0 COMMENT '已重试次数',
  max_retry INT NOT NULL DEFAULT 5 COMMENT '最大重试次数',
  next_retry_time DATETIME COMMENT '下次重试时间',
  create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  KEY idx_status_next_retry (status, next_retry_time),
  KEY idx_biz_key (biz_key),
  KEY idx_create_time (create_time)
) COMMENT='本地消息表';

-- 幂等表(消费端)
CREATE TABLE idempotent_record (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  biz_key VARCHAR(64) NOT NULL COMMENT '业务唯一键',
  handler VARCHAR(64) NOT NULL COMMENT '处理器标识',
  status TINYINT NOT NULL DEFAULT 0 COMMENT '0=处理中 1=已完成',
  create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
  UNIQUE KEY uk_biz_handler (biz_key, handler)
) COMMENT='幂等处理记录表';

2. 订单服务:写入业务数据 + 消息表(同一事务)

java
@Service
@Slf4j
public class OrderService {

    @Autowired
    private OrderRepository orderRepository;
    @Autowired
    private LocalMessageRepository messageRepository;
    @Autowired
    private TransactionTemplate transactionTemplate;

    public void createOrder(OrderDTO dto) {
        transactionTemplate.execute(status -> {
            // 1. 创建订单
            Order order = new Order();
            order.setUserId(dto.getUserId());
            order.setAmount(dto.getAmount());
            order.setStatus(OrderStatus.CREATED);
            orderRepository.save(order);

            // 2. 写入消息表
            DeductStockMessage msg = new DeductStockMessage(
                order.getId(), dto.getProductId(), dto.getQuantity()
            );
            LocalMessage msgRecord = new LocalMessage();
            msgRecord.setMsgId(UUID.randomUUID().toString());
            msgRecord.setBizKey("deduct_stock_" + order.getId()); // 业务唯一键
            msgRecord.setContent(JSON.toJSONString(msg));
            msgRecord.setStatus(0);
            msgRecord.setNextRetryTime(LocalDateTime.now());
            messageRepository.save(msgRecord);

            log.info("订单创建成功, orderId={}, msgId={}", order.getId(), msgRecord.getMsgId());
            return null;
        });
    }
}

3. 定时任务:轮询 + 发送消息(带分表扫描)

java
@Component
@Slf4j
public class MessageSenderJob {

    @Autowired
    private LocalMessageRepository messageRepository;
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 1000) // 高频场景缩到 1 秒
    public void sendPendingMessages() {
        // 扫描 16 张分表,每张查 50 条,总量最多 800 条
        for (int shard = 0; shard < 16; shard++) {
            List<LocalMessage> pending = messageRepository
                .findTop50ByStatusAndNextRetryTimeBefore(shard, 0, LocalDateTime.now());

            for (LocalMessage msg : pending) {
                try {
                    // 乐观锁:CAS 更新状态,防止多节点重复发送
                    int updated = messageRepository.casUpdateStatus(
                        msg.getId(), 0, 1, LocalDateTime.now(), shard
                    );
                    if (updated == 0) {
                        continue; // 被其他节点抢走了
                    }

                    // 发送到 Kafka,设置超时 3 秒
                    SendResult result = kafkaTemplate.send(
                        "stock-deduct", msg.getMsgId(), msg.getContent()
                    ).get(3, TimeUnit.SECONDS);

                    // 发送成功,标记为"已发送"
                    messageRepository.updateStatusById(msg.getId(), 1, 2, shard);
                    log.info("消息发送成功, msgId={}, partition={}", msg.getMsgId(), result.getRecordMetadata().partition());

                } catch (TimeoutException e) {
                    log.warn("Kafka 发送超时, msgId={}, retry={}", msg.getMsgId(), msg.getRetryCount());
                    // 超时回滚状态,下次重试
                    messageRepository.updateStatusById(msg.getId(), 1, 0, shard);

                } catch (Exception e) {
                    log.warn("消息发送失败, msgId={}, retry={}", msg.getMsgId(), msg.getRetryCount(), e);
                    // 指数退避:1,2,4,8,16 秒
                    long backoff = (long) Math.pow(2, Math.min(msg.getRetryCount(), 5));
                    messageRepository.updateRetry(msg.getId(),
                        msg.getRetryCount() + 1,
                        LocalDateTime.now().plusSeconds(backoff),
                        shard
                    );
                    // 超过最大重试,标记为死信
                    if (msg.getRetryCount() + 1 >= msg.getMaxRetry()) {
                        messageRepository.updateStatusById(msg.getId(), 0, 3, shard);
                        log.warn("消息达到最大重试次数, msgId={}, 转入死信", msg.getMsgId());
                    }
                }
            }
        }
    }
}

4. 消费端:幂等处理(唯一键 + 业务补偿)

java
@Component
@Slf4j
public class StockConsumer {

    @Autowired
    private StockService stockService;
    @Autowired
    private IdempotentRepository idempotentRepository;

    @KafkaListener(topics = "stock-deduct", groupId = "stock-group", concurrency = "3")
    public void onMessage(ConsumerRecord<String, String> record) {
        DeductStockMessage msg = JSON.parseObject(record.value(), DeductStockMessage.class);
        String bizKey = "deduct_stock_" + msg.getOrderId();

        // 幂等检查:用 UNIQUE KEY 做去重
        try {
            idempotentRepository.insert(bizKey, "stock_deduct", 0);
        } catch (DuplicateKeyException e) {
            log.info("重复消息跳过, bizKey={}", bizKey);
            return;
        }

        try {
            // 扣库存
            int result = stockService.deduct(msg.getProductId(), msg.getQuantity());
            if (result == 0) {
                // 库存不足,转人工处理
                log.warn("库存不足, productId={}, quantity={}", msg.getProductId(), msg.getQuantity());
                // 触发补偿流程:订单取消 + 退款
                compensationService.triggerOrderCancel(msg.getOrderId(), "库存不足");
            }
            // 更新幂等记录为已完成
            idempotentRepository.updateStatus(bizKey, "stock_deduct", 0, 1);
            log.info("库存扣减成功, orderId={}, productId={}, quantity={}",
                msg.getOrderId(), msg.getProductId(), msg.getQuantity());

        } catch (Exception e) {
            log.error("库存扣减失败, msgId={}", msg.getMsgId(), e);
            // 删除幂等记录,让重投能重新进来
            idempotentRepository.delete(bizKey, "stock_deduct");
            // 不 ACK,Kafka 自动重投
            throw new RuntimeException(e);
        }
    }
}

本地消息表的致命缺陷

  1. 数据库写入放大:每笔业务写操作附带至少一条消息表写入,高峰时 WAL 写入量翻倍
  2. 定时任务延迟:最短 1 秒扫描间隔,对于要求毫秒级一致性的场景不够
  3. 死信积压:重试超限后转入死信,需要人工处理,半夜被报警叫醒
  4. 分库分表复杂度:消息表不拆分,单表几百万行后索引性能急剧下降

方案二:事务消息(Transactional Message,RocketMQ)

RocketMQ 把"发消息"拆成两阶段,与本地事务绑定:

  1. 生产者发送半消息(half message)——消息到达 Broker,但消费者不可见
  2. 生产者执行本地事务
  3. 本地事务成功 → 提交半消息,消费者可见;本地事务失败 → 回滚半消息,消费者永远不会看到

如果第 2 步执行过程中生产者宕机,RocketMQ Broker 会定时回调生产者的 check 接口,询问这条半消息对应的事务到底是什么结果。生产者根据本地事务状态回复 commit 或 rollback。

事务消息完整实现

java
// 生产者配置
@Configuration
public class TransactionProducerConfig {

    @Bean
    public TransactionMQProducer transactionProducer() {
        TransactionMQProducer producer = new TransactionMQProducer("order-group");
        producer.setNamesrvAddr("192.168.1.100:9876");

        // 设置线程池,避免 check 回调阻塞主线程
        ExecutorService executor = new ThreadPoolExecutor(
            2, 4, 100, TimeUnit.SECONDS,
            new ArrayBlockingQueue<>(2000),
            new ThreadFactoryBuilder().setNameFormat("tx-check-pool-%d").build()
        );
        producer.setExecutorService(executor);

        // 设置事务监听器
        producer.setTransactionListener(new TransactionListener() {
            @Override
            public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
                OrderDTO dto = (OrderDTO) arg;
                try {
                    // 执行本地事务:创建订单
                    createOrderInLocalDb(dto);
                    return LocalTransactionState.COMMIT_MESSAGE;
                } catch (DuplicateKeyException e) {
                    // 幂等:订单已存在,按 commit 处理
                    log.warn("订单已存在, orderId={}, 按commit处理", dto.getOrderId());
                    return LocalTransactionState.COMMIT_MESSAGE;
                } catch (Exception e) {
                    log.error("本地事务执行失败, orderId={}", dto.getOrderId(), e);
                    return LocalTransactionState.ROLLBACK_MESSAGE;
                }
            }

            @Override
            public LocalTransactionState checkLocalTransaction(MessageExt msg) {
                // Broker 回调查询事务状态
                String orderId = msg.getKeys();
                boolean exists = orderService.existsById(orderId);
                if (exists) {
                    return LocalTransactionState.COMMIT_MESSAGE;
                }
                // 如果本地事务还在执行中,返回 UNKNOWN,Broker 会等下次回调
                if (isTransactionInProgress(orderId)) {
                    return LocalTransactionState.UNKNOW;
                }
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        });

        try {
            producer.start();
        } catch (MQClientException e) {
            log.error("事务消息生产者启动失败", e);
            throw new RuntimeException(e);
        }
        return producer;
    }

    private boolean isTransactionInProgress(String orderId) {
        // 从 Redis 或本地缓存查看事务执行状态
        return redisTemplate.hasKey("tx_progress:" + orderId);
    }
}

// 发送半消息
@Service
public class OrderService {

    @Autowired
    private TransactionMQProducer producer;

    public void createOrderWithTransaction(OrderDTO dto) {
        // 发送半消息前,先在 Redis 标记事务进行中
        redisTemplate.opsForValue().set("tx_progress:" + dto.getOrderId(), "1", 30, TimeUnit.SECONDS);

        try {
            Message msg = new Message("stock-deduct", "orderTag",
                dto.getOrderId().getBytes(StandardCharsets.UTF_8)
            );
            msg.setKeys(dto.getOrderId());

            SendResult result = producer.sendMessageInTransaction(msg, dto);
            log.info("事务消息发送结果: orderId={}, sendStatus={}", dto.getOrderId(), result.getSendStatus());

            if (result.getSendStatus() == SendStatus.SEND_OK) {
                // 事务提交成功,清除进度标记
                redisTemplate.delete("tx_progress:" + dto.getOrderId());
            }

        } catch (MQClientException e) {
            log.error("事务消息发送失败, orderId={}", dto.getOrderId(), e);
            throw new BizException("ORDER_CREATE_FAILED", "订单创建失败");
        }
    }
}

事务消息的坑:check 回调超时导致数据不一致

我遇到过一个线上 case:

  • 本地事务执行了 15 秒(因为订单关联了外部风控接口超时)
  • RocketMQ Broker 默认 6 秒后触发 check 回调
  • check 回调时本地事务还没完成,查数据库发现订单不存在,返回 ROLLBACK
  • 3 秒后本地事务执行成功,但半消息已经被回滚了
  • 订单创建了,库存没扣,变成脏数据

解决方案:半消息设置事务超时时间,或者保证 check 接口的幂等性:check 时如果发现本地事务还在执行中,返回 UNKNOWN 而不是 ROLLBACK,给 Broker 下次再查的机会。

java
// 生产端设置半消息超时
producer.setTransactionTimeOut(30); // 30 秒,给本地事务充足时间
producer.setCheckForbiddenTime(60); // 60 秒内 check 失败不报错

方案三:事件溯源(Event Sourcing)

不保存当前状态,只保存所有状态变更事件。当前状态由事件回放计算得出。

订单状态不是"订单表里一行 UPDATE 来 UPDATE 去",而是存了一串事件:OrderCreatedOrderPaidOrderShippedOrderDelivered。要查当前订单状态?把属于这个订单的所有事件按时间顺序回放一遍就知道。

事件溯源实现

java
// 事件基类
@Getter
public abstract class DomainEvent {
    private final String eventId = UUID.randomUUID().toString();
    private final String aggregateId;
    private final LocalDateTime occurredAt = LocalDateTime.now();
    private final int version;

    protected DomainEvent(String aggregateId, int version) {
        this.aggregateId = aggregateId;
        this.version = version;
    }
}

// 具体事件
public class OrderCreatedEvent extends DomainEvent {
    private final Long userId;
    private final BigDecimal amount;
    private final List<OrderItem> items;

    public OrderCreatedEvent(String orderId, int version, Long userId, BigDecimal amount, List<OrderItem> items) {
        super(orderId, version);
        this.userId = userId;
        this.amount = amount;
        this.items = items;
    }
}

// 事件仓库
@Repository
public class EventStore {

    @Autowired
    private JdbcTemplate jdbcTemplate;

    public void append(DomainEvent event) {
        String sql = "INSERT INTO domain_events (aggregate_id, event_type, version, event_data, occurred_at) " +
                     "VALUES (?, ?, ?, ?, ?)";
        jdbcTemplate.update(sql,
            event.getAggregateId(),
            event.getClass().getSimpleName(),
            event.getVersion(),
            JSON.toJSONString(event),
            event.getOccurredAt()
        );
    }

    public List<DomainEvent> loadEvents(String aggregateId) {
        String sql = "SELECT * FROM domain_events WHERE aggregate_id = ? ORDER BY version ASC";
        return jdbcTemplate.query(sql, new Object[]{aggregateId}, (rs, rowNum) -> {
            String eventType = rs.getString("event_type");
            String eventData = rs.getString("event_data");
            return JSON.parseObject(eventData, Class.forName(eventType));
        });
    }
}

// 聚合:通过事件回放重建状态
public class OrderAggregate {

    private String orderId;
    private OrderStatus status;
    private BigDecimal amount;
    private List<DomainEvent> changes = new ArrayList<>();

    // 从事件流重建
    public static OrderAggregate loadFromHistory(List<DomainEvent> events) {
        OrderAggregate aggregate = new OrderAggregate();
        for (DomainEvent event : events) {
            aggregate.apply(event);
        }
        return aggregate;
    }

    // 业务方法
    public void createOrder(Long userId, BigDecimal amount, List<OrderItem> items) {
        // 业务校验
        if (amount.compareTo(BigDecimal.ZERO) <= 0) {
            throw new IllegalArgumentException("金额必须大于0");
        }
        // 产生事件,不持久化
        int newVersion = changes.size() + 1;
        changes.add(new OrderCreatedEvent(orderId, newVersion, userId, amount, items));
    }

    // 应用事件
    private void apply(DomainEvent event) {
        if (event instanceof OrderCreatedEvent) {
            OrderCreatedEvent e = (OrderCreatedEvent) event;
            this.orderId = e.getAggregateId();
            this.status = OrderStatus.CREATED;
            this.amount = e.getAmount();
        }
        // 其他事件...
        this.version = event.getVersion();
    }

    public void save(EventStore store) {
        for (DomainEvent event : changes) {
            store.append(event);
        }
        changes.clear();
    }
}

事件溯源的实际案例:金融交易流水

我参与过一个金融交易系统,每天处理 2000 万笔交易,要求:

  • 任意一笔交易可追溯 3 年内的完整变更历史
  • 支持按时间点回放("2026-01-01 当时的账户余额是多少")
  • 审计要求:不能修改历史数据,只能补偿

用事件溯源 + CQRS(命令查询职责分离):

  • 写入端:EventStore 只追加,每秒 5000 事件写入
  • 查询端:定期生成快照(Snapshot),每 100 个事件打一个快照,查询时从最近快照开始回放
  • 快照表:snapshot(aggregate_id, version, state_json, created_at)
sql
-- 事件表
CREATE TABLE domain_events (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  aggregate_id VARCHAR(64) NOT NULL,
  event_type VARCHAR(64) NOT NULL,
  version INT NOT NULL,
  event_data JSON NOT NULL,
  occurred_at DATETIME(3) NOT NULL,
  UNIQUE KEY uk_agg_version (aggregate_id, version),
  KEY idx_agg_occurred (aggregate_id, occurred_at)
) ENGINE=InnoDB;

-- 快照表
CREATE TABLE aggregate_snapshot (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  aggregate_id VARCHAR(64) NOT NULL,
  version INT NOT NULL,
  state_json JSON NOT NULL,
  created_at DATETIME(3) NOT NULL,
  UNIQUE KEY uk_agg_version (aggregate_id, version)
) ENGINE=InnoDB;

性能实测:

  • 单聚合 1000 个事件,从快照回放:< 5ms
  • 单聚合 1000 个事件,无快照全量回放:~50ms
  • 单聚合 10000 个事件,无快照:~500ms(开始出现性能问题)

结论:事件溯源必须配合快照,否则查询性能随事件量线性劣化。

三大方案对比表

维度本地消息表事务消息 (RocketMQ)事件溯源
核心思想业务+消息同本地事务半消息+本地事务回调保存事件而非状态
MQ 依赖任意 MQ(Kafka/RabbitMQ/RocketMQ)仅 RocketMQ任意 MQ 或事件总线
维护成本高:需要消息表、定时任务、死信处理中:需实现 check 回调接口高:需要事件存储、快照、CQRS
延迟1-5 秒(定时任务扫描间隔)毫秒级(半消息提交后立即可见)毫秒级(事件总线发布)
查询复杂度低:直接 SQL 查表低:同本地消息表高:需事件回放或维护物化视图
审计日志不天然支持不天然支持天然支持,所有变更可追溯
统一业务场景订单、支付、库存(通用场景)强一致性消息传递场景财务流水、审计日志、状态机
代码侵入中:需额外消息表操作中:需实现 TransactionListener高:事件驱动重构,非 CRUD 思维
幂等性要求必须必须天然幂等(事件幂等追加)
分库分表支持需要分表不需要(MQ 管理)按 aggregate_id 分区

什么时候选哪个?

  1. 团队用 Kafka/RabbitMQ,且能接受运维负担 → 本地消息表。但必须做分表,且定时任务独立线程池。

  2. 团队已经在用 RocketMQ 4.x+ → 事务消息。运维成本最低,延迟最低。但注意 check 回调超时问题。

  3. 业务有强制审计要求,或状态机非常复杂 → 事件溯源。但要做好心理准备:团队需要理解 CQRS 模式,查询端需要额外维护物化视图。

  4. 业务量极低(日均几百单) → 直接本地消息表就行,不需要分表,不需要 RocketMQ。

所有方案公用的底线

幂等性不是可选项,是必须的。 无论哪种方案,消费端必须用业务唯一键做去重。某电商公司双十一因为幂等表没加 UNIQUE KEY,重复消息导致多扣了 2 万件库存,凌晨 3 点被 DBA 叫起来做数据订正。

幂等实现三要素:

  1. 业务唯一键(如 orderId + 业务类型)
  2. 唯一约束或分布式锁保证先到先得
  3. 消费失败时删除幂等记录,允许重试重新进入

另外,所有最终一致性方案都有一个共同代价:数据不一致的时间窗口。如果业务不能接受 1 秒以上的不一致(比如银行转账),那就不该用微服务,或者用 TCC 补偿事务 + 冻结资金的模式。

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。