微服务架构下,跨服务的数据操作无法依赖数据库本地事务保证ACID特性。分布式事务最终一致性方案通过消息中间件实现跨服务数据同步,是高并发设计中的核心组件。本文对比本地消息表和事务消息两种方案,给出选型建议和代码实现。
分布式事务问题场景分析
典型场景:电商系统下单后需要同时操作订单服务(创建订单)和库存服务(扣减库存)。两个服务使用独立数据库,本地事务无法跨库。如果订单创建成功但库存扣减失败,会出现超卖问题。
// 问题代码 - 非事务性操作
@Transactional
public void createOrder(OrderDTO dto) {
// 1. 本地:创建订单
orderMapper.insert(dto);
// 2. 远程:扣减库存(可能失败)
inventoryFeignClient.deduct(dto.getProductId(), dto.getQuantity());
// 如果这里抛异常,订单已创建但库存未扣减
}
强一致性方案(如Seata AT模式)通过全局锁实现跨库事务,但性能损耗大,不适合高并发场景。最终一致性方案保证数据最终一致,在可接受短暂不一致的业务场景中是更优选择。
本地消息表方案实现
本地消息表的核心思路:将消息发送操作和业务操作放在同一个数据库事务中,保证业务操作和消息记录要么同时成功要么同时失败。后台任务轮询消息表,将未发送的消息投递到消息中间件。
-- 消息表DDL
CREATE TABLE local_message (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
business_key VARCHAR(64) NOT NULL COMMENT '业务唯一键',
topic VARCHAR(64) NOT NULL COMMENT '消息主题',
message_body TEXT NOT NULL COMMENT '消息内容JSON',
status TINYINT DEFAULT 0 COMMENT '0-待发送 1-已发送 2-失败',
retry_count INT DEFAULT 0 COMMENT '重试次数',
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_status_create (status, create_time)
) ENGINE=InnoDB;
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private MessageMapper messageMapper;
@Autowired
private RocketMQTemplate mqTemplate;
@Transactional(rollbackFor = Exception.class)
public void createOrder(OrderDTO dto) {
// 1. 创建订单
Order order = new Order();
order.setOrderId(dto.getOrderId());
order.setProductId(dto.getProductId());
order.setQuantity(dto.getQuantity());
order.setStatus("CREATED");
orderMapper.insert(order);
// 2. 同事务写入消息表
LocalMessage msg = new LocalMessage();
msg.setBusinessKey(dto.getOrderId());
msg.setTopic("inventory-deduct");
msg.setMessageBody(JSON.toJSONString(dto));
msg.setStatus(0);
messageMapper.insert(msg);
// 事务提交后,订单和消息记录同时落库
}
}
// 消息投递定时任务
@Component
public class MessageSendTask {
@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
List<LocalMessage> messages = messageMapper.selectPending(100);
for (LocalMessage msg : messages) {
try {
mqTemplate.convertAndSend(msg.getTopic(), msg.getMessageBody());
messageMapper.updateStatus(msg.getId(), 1);
} catch (Exception e) {
int retry = msg.getRetryCount() + 1;
if (retry >= 5) {
messageMapper.updateStatus(msg.getId(), 2);
// 告警通知
} else {
messageMapper.updateRetry(msg.getId(), retry);
}
}
}
}
}
本地消息表方案的优势:实现简单,不依赖消息中间件的事务特性,可适配任何MQ。劣势:定时轮询有延迟(通常5-10秒),消息表数据持续增长需要定期归档。
RocketMQ事务消息方案实现
事务消息方案依赖RocketMQ的两阶段发送机制。先发送半消息(对消费者不可见),执行本地事务后根据结果提交或回滚半消息。如果本地事务执行超时,RocketMQ回查本地事务状态。
@Service
public class OrderTransactionService {
@Autowired
private RocketMQTemplate mqTemplate;
@Autowired
private OrderMapper orderMapper;
public void createOrder(OrderDTO dto) {
Message<String> message = MessageBuilder
.withPayload(JSON.toJSONString(dto))
.setHeader("businessKey", dto.getOrderId())
.build();
mqTemplate.sendMessageInTransaction(
"inventory-deduct",
message,
dto
);
}
@RocketMQTransactionListener
public class OrderTransactionListener
implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
OrderDTO dto = (OrderDTO) arg;
try {
Order order = new Order();
order.setOrderId(dto.getOrderId());
order.setProductId(dto.getProductId());
order.setQuantity(dto.getQuantity());
order.setStatus("CREATED");
orderMapper.insert(order);
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderId = (String) msg.getHeaders().get("businessKey");
Order order = orderMapper.selectByOrderId(orderId);
if (order != null) {
return RocketMQLocalTransactionState.COMMIT;
}
return RocketMQLocalTransactionState.ROLLBACK;
}
}
}
事务消息方案的优势:消息投递实时性好(毫秒级),无需轮询,RocketMQ自带事务回查机制保证可靠性。劣势:强依赖RocketMQ,无法平滑切换到其他MQ。服务治理中需要额外维护事务回查接口。
两种方案选型对比与生产建议
选型决策矩阵:
维度 本地消息表 事务消息
消息延迟 5-10秒 毫秒级
MQ兼容性 任意MQ 仅RocketMQ
实现复杂度 中等 较高
数据库压力 有额外写 无额外表
事务回查 不支持 支持
运维成本 需归档清理 低
业务中台建设的实践建议:对消息实时性要求高的核心链路(支付、库存)使用事务消息方案;对实时性要求不高的旁路操作(通知、日志同步)使用本地消息表方案。API接口规范中应定义消息消费的幂等性要求,消费端通过business_key做去重。
// 消费端幂等处理
@RocketMQMessageListener(topic = "inventory-deduct")
public class InventoryConsumer implements RocketMQListener<OrderDTO> {
@Autowired
private InventoryMapper inventoryMapper;
@Autowired
private RedisTemplate redisTemplate;
@Override
@Transactional
public void onMessage(OrderDTO dto) {
String key = "deduct:" + dto.getOrderId();
Boolean acquired = redisTemplate.opsForValue()
.setIfAbsent(key, "1", 24, TimeUnit.HOURS);
if (Boolean.FALSE.equals(acquired)) {
return; // 已处理,跳过
}
int rows = inventoryMapper.deduct(
dto.getProductId(), dto.getQuantity());
if (rows == 0) {
redisTemplate.delete(key); // 回滚幂等标记
throw new RuntimeException("库存不足");
}
}
}
幂等性是最终一致性方案的必要补充。消息中间件保证at-least-once投递语义,消费端必须处理重复消息。Redis SETNX配合业务唯一键是轻量级幂等方案,在高并发场景下性能表现稳定。
分布式事务最终一致性方案的选择没有银弹。理解业务对一致性和实时性的要求,结合团队能力和技术栈选型,才能做出合理的架构决策。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-zui-zhong-yi-zhi-xing-fang-an-ben-di-xiao/