MQ事务消息详解:RocketMQ的分布式事务方案
·
MQ事务消息详解:RocketMQ的分布式事务方案 🔄
|
🌺The Begin🌺点点关注,收藏不迷路🌺
|
引言:从本地消息表到MQ事务消息
本地消息表方案虽然简单可靠,但有一个明显缺点:消息表占用业务数据库资源。
MQ事务消息本质上是对本地消息表的封装,整体流程与本地消息表一致,唯一不同的就是将本地消息表存在了MQ内部,而不是业务数据库中。
一、MQ事务消息的核心思想 🎯
1.1 本质:消息表上移到MQ
| 对比维度 | 本地消息表 | MQ事务消息 |
|---|---|---|
| 消息存储 | 业务数据库 | MQ内部 |
| 事务保证 | 本地数据库事务 | 半消息机制 + 回查 |
| 资源消耗 | 占用业务库资源 | 不占用业务库 |
| 实现复杂度 | 简单 | 中等 |
| 适用MQ | 任何MQ | 需支持事务消息 |
1.2 RocketMQ事务消息的演进
- RocketMQ 4.3之前:不支持事务消息,需借助本地消息表
- RocketMQ 4.3之后:内置事务消息支持,完美封装了本地消息表的思想
二、RocketMQ事务消息执行流程 📋
2.1 整体流程
2.2 三个核心角色
| 角色 | 职责 | 类比 |
|---|---|---|
| 事务发起方 | 执行本地事务,提交/回滚消息 | 订单服务 |
| RocketMQ | 存储半消息,回查事务状态 | 消息协调器 |
| 事务被动方 | 消费消息,执行业务 | 库存服务 |
三、详细执行步骤 🔍
3.1 第一阶段:发送半消息
// 事务发起方代码
@Service
public class OrderService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void createOrder(OrderDTO dto) {
// 1. 构建消息
Message message = MessageBuilder.withPayload(dto)
.setHeader("orderId", dto.getOrderId())
.build();
// 2. 发送半消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-topic", // topic
message, // 消息内容
dto // 本地事务需要的参数
);
if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
log.info("事务提交成功");
} else {
log.error("事务失败");
}
}
}
半消息的特点:
- ✅ 对消费者不可见
- ✅ 存储在MQ中,不占用业务数据库
- ✅ 需要等待commit后才投递
3.2 第二阶段:执行本地事务
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderRepository orderRepository;
@Autowired
private TransactionLogRepository logRepository;
/**
* 执行本地事务
*/
@Override
@Transactional
public RocketMQLocalTransactionState executeLocalTransaction(
Message msg, Object arg) {
OrderDTO dto = (OrderDTO) arg;
String transactionId = msg.getHeaders().get("rocketmq_TRANSACTION_ID").toString();
try {
// 1. 记录事务日志(用于回查)
TransactionLog log = new TransactionLog();
log.setTransactionId(transactionId);
log.setBusinessKey(dto.getOrderId());
log.setStatus(0); // 执行中
logRepository.save(log);
// 2. 执行业务操作
Order order = new Order();
order.setId(dto.getOrderId());
order.setUserId(dto.getUserId());
order.setAmount(dto.getAmount());
order.setStatus("PENDING");
orderRepository.save(order);
// 3. 更新事务日志
log.setStatus(1); // 成功
logRepository.save(log);
// 4. 返回提交
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
log.error("本地事务失败", e);
// 更新事务日志状态为失败
TransactionLog log = logRepository.findByTransactionId(transactionId);
if (log != null) {
log.setStatus(2); // 失败
logRepository.save(log);
}
return RocketMQLocalTransactionState.ROLLBACK;
}
}
}
3.3 回查机制
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private TransactionLogRepository logRepository;
/**
* 回查事务状态
* RocketMQ会定时调用此方法检查未决事务
*/
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String transactionId = msg.getHeaders().get("rocketmq_TRANSACTION_ID").toString();
// 查询事务日志
TransactionLog log = logRepository.findByTransactionId(transactionId);
if (log == null) {
// 没有日志,说明事务未执行,回滚
return RocketMQLocalTransactionState.ROLLBACK;
}
// 根据日志状态返回
switch (log.getStatus()) {
case 1: // 成功
return RocketMQLocalTransactionState.COMMIT;
case 2: // 失败
return RocketMQLocalTransactionState.ROLLBACK;
default: // 执行中,等待下次回查
return RocketMQLocalTransactionState.UNKNOWN;
}
}
}
3.4 消费端实现
@Service
@RocketMQMessageListener(
topic = "order-topic",
consumerGroup = "stock-consumer-group"
)
public class StockConsumer implements RocketMQListener<OrderDTO> {
@Autowired
private StockService stockService;
@Autowired
private ConsumeLogService consumeLogService;
@Override
public void onMessage(OrderDTO dto) {
String messageId = dto.getMessageId();
// 1. 幂等检查
if (consumeLogService.isProcessed(messageId)) {
log.info("消息已处理,跳过:{}", messageId);
return;
}
// 2. 执行业务
try {
stockService.deduct(dto.getProductId(), dto.getQuantity());
// 3. 记录消费日志
consumeLogService.save(messageId);
} catch (Exception e) {
log.error("库存扣减失败", e);
// 抛出异常,RocketMQ会重试
throw new RuntimeException(e);
}
}
}
四、事务日志表设计 📊
虽然消息存储在MQ,但回查仍需要事务日志:
CREATE TABLE `transaction_log` (
`id` BIGINT PRIMARY KEY AUTO_INCREMENT,
`transaction_id` VARCHAR(64) NOT NULL COMMENT '事务ID',
`business_key` VARCHAR(100) NOT NULL COMMENT '业务主键',
`status` TINYINT NOT NULL DEFAULT 0 COMMENT '状态:0-执行中 1-成功 2-失败',
`retry_count` INT DEFAULT 0 COMMENT '回查次数',
`create_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
`update_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE INDEX uk_transaction_id (`transaction_id`),
INDEX idx_status (`status`, `retry_count`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
五、事务状态流转 🔄
六、容错处理机制 🛡️
6.1 各种异常场景的处理
| 异常场景 | 处理方式 | 最终结果 |
|---|---|---|
| 发送半消息失败 | 业务直接失败 | ✅ 数据一致 |
| 本地事务成功,提交失败 | MQ回查,提交消息 | ✅ 数据一致 |
| 本地事务成功,提交超时 | MQ回查,提交消息 | ✅ 数据一致 |
| 本地事务失败,回滚成功 | 半消息删除 | ✅ 数据一致 |
| 本地事务失败,回滚失败 | MQ回查,回滚消息 | ✅ 数据一致 |
| 事务状态未知 | MQ定时回查 | ✅ 最终一致 |
6.2 幂等保证
public class IdempotentConsumer {
// 使用业务ID做幂等
private Set<String> processedMessages =
Collections.newSetFromMap(new ConcurrentHashMap<>());
public void consume(String messageId, BusinessData data) {
// 1. 去重
if (!processedMessages.add(messageId)) {
log.info("消息重复消费,已过滤");
return;
}
// 2. 执行业务
try {
doBusiness(data);
// 3. 持久化消息ID(用于重启后去重)
saveToDB(messageId);
} catch (Exception e) {
processedMessages.remove(messageId);
throw e;
}
}
}
七、本地消息表 vs RocketMQ事务消息 📊
| 对比维度 | 本地消息表 | RocketMQ事务消息 |
|---|---|---|
| 消息存储 | 业务数据库 | MQ内部 |
| 资源占用 | 占用业务库 | 不占用业务库 |
| 实现复杂度 | 简单 | 中等 |
| 性能 | 受限于数据库 | 高(MQ专用存储) |
| 适用MQ | 任何MQ | 仅RocketMQ 4.3+ |
| 回查机制 | 定时任务扫描 | MQ内置回查 |
| 事务日志 | 完整消息表 | 精简事务日志 |
八、最佳实践建议 💡
8.1 事务日志清理
@Component
public class TransactionLogCleaner {
@Scheduled(cron = "0 0 2 * * ?") // 每天凌晨2点
public void cleanOldLogs() {
// 删除7天前的日志
LocalDateTime deadline = LocalDateTime.now().minusDays(7);
logRepository.deleteByCreateTimeBefore(deadline);
}
}
8.2 回查次数控制
public class CheckLocalTransaction {
private static final int MAX_CHECK_COUNT = 15;
public RocketMQLocalTransactionState check(Message msg) {
String transactionId = getTransactionId(msg);
TransactionLog log = logRepository.findByTransactionId(transactionId);
if (log.getRetryCount() >= MAX_CHECK_COUNT) {
// 超过最大回查次数,标记为失败
return RocketMQLocalTransactionState.ROLLBACK;
}
// 更新回查次数
log.setRetryCount(log.getRetryCount() + 1);
logRepository.save(log);
// 根据状态判断
return determineState(log);
}
}
8.3 监控告警
@Component
public class TransactionMonitor {
@Scheduled(fixedDelay = 60000)
public void monitor() {
// 查询长时间未决的事务
List<TransactionLog> pendingLogs = logRepository
.findByStatusAndCreateTimeBefore(0, LocalDateTime.now().minusMinutes(10));
if (!pendingLogs.isEmpty()) {
// 发送告警
alertService.send("发现长时间未决事务", pendingLogs);
}
}
}
九、总结 🎯
9.1 一句话总结
RocketMQ事务消息是对本地消息表的完美封装,将消息存储从业务数据库迁移到MQ内部,既保留了本地消息表的可靠性,又解决了资源占用问题,是分布式事务的优雅实现。
9.2 核心价值
| 价值 | 说明 |
|---|---|
| 资源隔离 | 消息不占用业务库资源 |
| 高性能 | MQ专用存储,高吞吐 |
| 自动化 | 内置回查机制 |
| 可靠 | 不丢失消息,最终一致 |
9.3 适用场景
public class Advice {
public String getAdvice() {
return "1. 新项目首选RocketMQ事务消息\n"
+ "2. 高并发场景推荐使用\n"
+ "3. 记得做好幂等和日志清理\n"
+ "4. 监控未决事务,及时处理";
}
}
(本文为分布式系统事务系列文章,欢迎关注更多架构深度内容)

|
🌺The End🌺点点关注,收藏不迷路🌺
|
更多推荐

所有评论(0)