🌺The Begin🌺点点关注,收藏不迷路🌺

引言:从本地消息表到MQ事务消息

本地消息表方案虽然简单可靠,但有一个明显缺点:消息表占用业务数据库资源

MQ事务消息本质上是对本地消息表的封装,整体流程与本地消息表一致,唯一不同的就是将本地消息表存在了MQ内部,而不是业务数据库中。

MQ事务消息

MQ

消息存储

事务状态表

本地消息表

业务库

业务表

消息表


一、MQ事务消息的核心思想 🎯

1.1 本质:消息表上移到MQ

对比维度本地消息表MQ事务消息
消息存储业务数据库MQ内部
事务保证本地数据库事务半消息机制 + 回查
资源消耗占用业务库资源不占用业务库
实现复杂度简单中等
适用MQ任何MQ需支持事务消息

1.2 RocketMQ事务消息的演进

  • RocketMQ 4.3之前:不支持事务消息,需借助本地消息表
  • RocketMQ 4.3之后:内置事务消息支持,完美封装了本地消息表的思想

二、RocketMQ事务消息执行流程 📋

2.1 整体流程

事务被动方 业务数据库 RocketMQ 事务发起方 事务被动方 业务数据库 RocketMQ 事务发起方 第一阶段:发送半消息 第二阶段:执行本地事务 alt [业务成功] [业务失败] [本地事务超时或异常] 1. 发送半消息(prepare) 半消息保存成功 2. 执行业务操作 业务执行结果 3. 提交消息(commit) 4. 投递消息 5. 消费业务 6. 消费成功ACK 3. 回滚消息(rollback) 4. 删除半消息 3. 回查事务状态 4. 返回事务结果

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;

五、事务状态流转 🔄

发送prepare

执行本地事务

COMMIT

ROLLBACK

UNKNOWN

定时任务

事务成功

事务失败

待下次回查

投递消息

删除半消息

半消息

本地事务

提交

回滚

未知

回查

消息可见

消息删除


六、容错处理机制 🛡️

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🌺点点关注,收藏不迷路🌺
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐