微服务数据一致性:最终一致性 + 本地消息表 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,消费者消费成功后回调确认。
关键链路:
- 本地事务保证了业务数据和消息的原子性
- 定时任务保证了消息"至少一次"的投递
- 消费端幂等表保证了"恰好一次"的处理
实际踩坑:某电商双十一消息表写爆
我参与过一个电商项目,双十一当天本地消息表每秒写入 8000+ 条消息,定时任务每 5 秒扫一次,一次查 100 条。结果:
- 消息表膨胀到 500 万行,
idx_status_next_retry索引失效,扫一次 3 秒 - 订单都处理完了,库存还没扣,用户付了钱但发不出货
- 尾段时间把扫描间隔缩到 1 秒,一次查 500 条,数据库连接池直接打满
教训:本地消息表方案必须做分表 + 单独的扫描线程池,不能和业务 API 共享连接池。
代码实现
1. 消息表设计(带分表键)
-- 按 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. 订单服务:写入业务数据 + 消息表(同一事务)
@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. 定时任务:轮询 + 发送消息(带分表扫描)
@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. 消费端:幂等处理(唯一键 + 业务补偿)
@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);
}
}
}本地消息表的致命缺陷
- 数据库写入放大:每笔业务写操作附带至少一条消息表写入,高峰时 WAL 写入量翻倍
- 定时任务延迟:最短 1 秒扫描间隔,对于要求毫秒级一致性的场景不够
- 死信积压:重试超限后转入死信,需要人工处理,半夜被报警叫醒
- 分库分表复杂度:消息表不拆分,单表几百万行后索引性能急剧下降
方案二:事务消息(Transactional Message,RocketMQ)
RocketMQ 把"发消息"拆成两阶段,与本地事务绑定:
- 生产者发送半消息(half message)——消息到达 Broker,但消费者不可见
- 生产者执行本地事务
- 本地事务成功 → 提交半消息,消费者可见;本地事务失败 → 回滚半消息,消费者永远不会看到
如果第 2 步执行过程中生产者宕机,RocketMQ Broker 会定时回调生产者的 check 接口,询问这条半消息对应的事务到底是什么结果。生产者根据本地事务状态回复 commit 或 rollback。
事务消息完整实现
// 生产者配置
@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 下次再查的机会。
// 生产端设置半消息超时
producer.setTransactionTimeOut(30); // 30 秒,给本地事务充足时间
producer.setCheckForbiddenTime(60); // 60 秒内 check 失败不报错方案三:事件溯源(Event Sourcing)
不保存当前状态,只保存所有状态变更事件。当前状态由事件回放计算得出。
订单状态不是"订单表里一行 UPDATE 来 UPDATE 去",而是存了一串事件:OrderCreated → OrderPaid → OrderShipped → OrderDelivered。要查当前订单状态?把属于这个订单的所有事件按时间顺序回放一遍就知道。
事件溯源实现
// 事件基类
@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)
-- 事件表
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 分区 |
什么时候选哪个?
团队用 Kafka/RabbitMQ,且能接受运维负担 → 本地消息表。但必须做分表,且定时任务独立线程池。
团队已经在用 RocketMQ 4.x+ → 事务消息。运维成本最低,延迟最低。但注意 check 回调超时问题。
业务有强制审计要求,或状态机非常复杂 → 事件溯源。但要做好心理准备:团队需要理解 CQRS 模式,查询端需要额外维护物化视图。
业务量极低(日均几百单) → 直接本地消息表就行,不需要分表,不需要 RocketMQ。
所有方案公用的底线
幂等性不是可选项,是必须的。 无论哪种方案,消费端必须用业务唯一键做去重。某电商公司双十一因为幂等表没加 UNIQUE KEY,重复消息导致多扣了 2 万件库存,凌晨 3 点被 DBA 叫起来做数据订正。
幂等实现三要素:
- 业务唯一键(如 orderId + 业务类型)
- 唯一约束或分布式锁保证先到先得
- 消费失败时删除幂等记录,允许重试重新进入
另外,所有最终一致性方案都有一个共同代价:数据不一致的时间窗口。如果业务不能接受 1 秒以上的不一致(比如银行转账),那就不该用微服务,或者用 TCC 补偿事务 + 冻结资金的模式。