事务消息案例

wen java案例 2

本文目录导读:

事务消息案例

  1. 事务消息核心概念
  2. 订单-库存事务消息案例
  3. 完整代码实现
  4. 其他业务场景案例
  5. 事务消息最佳实践
  6. 事务消息最佳实践总结
  7. 生产环境注意事项

我将为您详细介绍RocketMQ事务消息的完整案例,涵盖分布式事务场景下的消息一致性解决方案。

事务消息核心概念

1 应用场景

  • 订单创建与库存扣减:创建订单后需要异步扣减库存
  • 支付确认与积分发放:支付成功后发放用户积分
  • 账户转账:跨账户资金转移
  • 电商下单流程:下单、锁库存、生成物流单

2 核心机制

发送半消息 → 执行本地事务 → 提交/回滚消息 → 服务端确认
     ↓              ↓              ↓
  Broker存储    本地事务执行   消费者可见
     ↓              ↓
  事务回查 ←—— 超时未决事务

订单-库存事务消息案例

1 业务流程图

sequenceDiagram
    participant App as 订单服务
    participant Broker as RocketMQ Broker
    participant DB as 订单数据库
    participant Inv as 库存服务
    App->>Broker: 发送半消息(order_tx_topic)
    Broker-->>App: 返回发送结果
    App->>DB: 执行本地事务(创建订单)
    alt 本地事务成功
        App->>Broker: 提交事务消息
        Broker->>Inv: 投递消息到库存服务
        Inv->>Inv: 扣减库存
    else 本地事务失败
        App->>Broker: 回滚事务消息
    end
    Note over App,Broker: 如果超时未决,Broker发起事务回查

完整代码实现

1 订单服务端(事务发送方)

@Component
public class OrderServiceImpl implements OrderService {
    @Autowired
    private TransactionMQProducer transactionProducer;
    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private RedisTemplate<String, String> redisTemplate;
    private static final String TX_TOPIC = "order-tx-topic";
    /**
     * 创建订单(事务消息核心方法)
     */
    public void createOrder(OrderDTO orderDTO) {
        String orderId = generateOrderId();
        orderDTO.setOrderId(orderId);
        // 1. 构建消息
        Message message = new Message(TX_TOPIC, 
                                     orderId.getBytes(),
                                     orderDTO);
        message.setKeys(orderId);
        // 2. 发送事务消息
        TransactionSendResult sendResult = transactionProducer.sendMessageInTransaction(
            message,
            new LocalTransactionExecutor() {
                @Override
                public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
                    return executeLocalOrderTransaction(orderDTO);
                }
                @Override
                public LocalTransactionState checkLocalTransaction(MessageExt msg) {
                    return checkOrderTransactionStatus(msg);
                }
            },
            orderDTO
        );
        // 3. 处理发送结果
        if (sendResult.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
            log.info("订单创建成功并提交事务消息,orderId: {}", orderId);
        } else {
            log.error("订单创建失败,orderId: {}", orderId);
            throw new BusinessException("订单创建失败");
        }
    }
    /**
     * 执行本地事务(创建订单记录)
     */
    private LocalTransactionState executeLocalOrderTransaction(OrderDTO orderDTO) {
        try {
            // 1. 创建本地订单
            Order order = new Order();
            BeanUtils.copyProperties(orderDTO, order);
            order.setStatus(OrderStatus.CREATED);
            order.setCreateTime(new Date());
            // 2. 插入订单表
            int insertCount = orderMapper.insert(order);
            if (insertCount > 0) {
                // 3. 保存订单状态到Redis(用于事务回查)
                String key = "tx_order_" + orderDTO.getOrderId();
                redisTemplate.opsForValue().set(key, 
                    JSON.toJSONString(order), 
                    30, TimeUnit.MINUTES);
                log.info("本地事务执行成功,订单创建完成");
                return LocalTransactionState.COMMIT_MESSAGE;
            } else {
                log.warn("本地事务执行失败,订单未创建");
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        } catch (Exception e) {
            log.error("本地事务执行异常", e);
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
    /**
     * 事务回查(解决消息发送中途挂掉)
     */
    private LocalTransactionState checkOrderTransactionStatus(MessageExt msg) {
        String orderId = msg.getKeys();
        log.info("收到事务回查请求,orderId: {}", orderId);
        // 1. 从Redis查询订单状态
        String orderJson = redisTemplate.opsForValue().get("tx_order_" + orderId);
        if (StringUtils.isNotEmpty(orderJson)) {
            Order order = JSON.parseObject(orderJson, Order.class);
            if (order.getStatus() == OrderStatus.CREATED) {
                log.info("订单存在且已创建,提交事务消息");
                return LocalTransactionState.COMMIT_MESSAGE;
            }
        }
        // 2. Redis没有,查询数据库(最终一致性保证)
        Order order = orderMapper.selectByOrderId(orderId);
        if (order != null) {
            log.info("数据库确认订单存在,提交事务");
            return LocalTransactionState.COMMIT_MESSAGE;
        }
        log.warn("订单不存在,回滚事务");
        return LocalTransactionState.ROLLBACK_MESSAGE;
    }
}

2 事务消息生产者配置

@Configuration
public class RocketMQConfig {
    @Bean
    public TransactionMQProducer transactionProducer() {
        TransactionMQProducer producer = new TransactionMQProducer("order-tx-producer-group");
        producer.setNamesrvAddr("127.0.0.1:9876");
        // 设置事务回查监听器
        producer.setTransactionListener(new TransactionListener() {
            @Override
            public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
                return null; // 在业务层实现
            }
            @Override
            public LocalTransactionState checkLocalTransaction(MessageExt msg) {
                return null; // 在业务层实现
            }
        });
        // 配置线程池
        ExecutorService executor = new ThreadPoolExecutor(
            5, 10, 60, TimeUnit.SECONDS,
            new ArrayBlockingQueue<>(100),
            new ThreadFactory() {
                private AtomicInteger count = new AtomicInteger();
                @Override
                public Thread newThread(Runnable r) {
                    return new Thread(r, "tx-msg-checker-" + count.incrementAndGet());
                }
            }
        );
        producer.setExecutorService(executor);
        // 启动
        try {
            producer.start();
        } catch (MQClientException e) {
            log.error("事务消息生产者启动失败", e);
        }
        return producer;
    }
}

3 库存服务端(消息消费者)

@Component
@RocketMQMessageListener(
    topic = "order-tx-topic",
    consumerGroup = "inventory-consumer-group",
    messageModel = MessageModel.CONCURRENTLY
)
public class InventoryConsumer implements RocketMQListener<MessageExt> {
    @Autowired
    private InventoryService inventoryService;
    @Override
    public void onMessage(MessageExt message) {
        // 1. 解析消息
        String body = new String(message.getBody());
        OrderDTO orderDTO = JSON.parseObject(body, OrderDTO.class);
        String topic = message.getTopic();
        String tags = message.getTags();
        String keys = message.getKeys();
        log.info("库存服务收到消息,订单号: {}, 商品ID: {}", 
                 keys, orderDTO.getOrderId());
        // 2. 幂等处理
        String dedupKey = String.format("dedup:inventory:%s:%s", 
                                        topic, keys);
        Boolean firstTime = redisTemplate.opsForValue()
            .setIfAbsent(dedupKey, "1", 10, TimeUnit.MINUTES);
        if (!Boolean.TRUE.equals(firstTime)) {
            log.info("消息已处理,跳过重复消费,key: {}", dedupKey);
            return;
        }
        try {
            // 3. 扣减库存
            inventoryService.deductStock(orderDTO);
            // 4. 记录消费成功
            log.info("库存扣减成功,订单号: {}", keys);
        } catch (Exception e) {
            log.error("库存扣减失败", e);
            // 抛出异常触发RocketMQ重试
            throw new RuntimeException("库存扣减失败", e);
        }
    }
}

其他业务场景案例

1 支付-积分事务

@Component
public class PaymentService {
    @Autowired
    private TransactionMQProducer txProducer;
    public void handlePayment(PaymentDTO paymentDTO) {
        Message message = new Message("payment-tx-topic",
                                     paymentDTO.getPaymentId().getBytes());
        TransactionSendResult result = txProducer.sendMessageInTransaction(
            message,
            (msg, arg) -> {
                // 本地事务:更新支付状态
                boolean success = paymentMapper.updateStatus(
                    paymentDTO.getPaymentId(), 
                    PaymentStatus.PAID
                );
                return success ? 
                    LocalTransactionState.COMMIT_MESSAGE : 
                    LocalTransactionState.ROLLBACK_MESSAGE;
            },
            paymentDTO
        );
    }
}
// 积分消费者
@Component
public class PointsConsumer {
    @RocketMQMessageListener(topic = "payment-tx-topic", 
                            consumerGroup = "points-group")
    public static class RewardsConsumer implements RocketMQListener<MessageExt> {
        @Override
        public void onMessage(MessageExt msg) {
            // 发放积分逻辑
            PaymentDTO paymentDTO = JSON.parseObject(
                new String(msg.getBody()), PaymentDTO.class
            );
            // 1. 计算积分
            int points = (int)(paymentDTO.getAmount() * 100);
            // 2. 发放积分
            pointsService.addPoints(paymentDTO.getUserId(), points);
            // 3. 记录日志
            log.info("用户 {} 获得积分 {}", 
                     paymentDTO.getUserId(), points);
        }
    }
}

2 转账-账户事务

public class TransferService {
    public void transfer(TransferDTO transferDTO) {
        // 生成转账流水号
        String transferId = UUID.randomUUID().toString();
        // 发送转账事务消息
        Message message = new Message("transfer-tx-topic",
                                      ("transfer_" + transferId).getBytes());
        TransactionSendResult result = producer.sendMessageInTransaction(
            message, 
            (msg, arg) -> {
                try {
                    // 本地事务:冻结转出账户资金
                    accountService.freezeAmount(transferDTO);
                    // 创建转账记录
                    transferMapper.insert(transferDTO);
                    return LocalTransactionState.COMMIT_MESSAGE;
                } catch (Exception e) {
                    return LocalTransactionState.ROLLBACK_MESSAGE;
                }
            },
            transferDTO
        );
    }
}

事务消息最佳实践

1 可靠性保障

public class ReliableMessageSender {
    public void sendReliableMessage(Message message) {
        // 1. 发送前持久化消息记录
        MessageRecord record = MessageRecord.builder()
            .msgId(message.getKeys())
            .topic(message.getTopic())
            .body(new String(message.getBody()))
            .status(MessageStatus.SENDING)
            .createTime(new Date())
            .build();
        messageRecordMapper.insert(record);
        // 2. 发送事务消息
        TransactionSendResult result = producer.sendMessageInTransaction(
            message, 
            (msg, arg) -> {
                // 3. 执行本地事务
                LocalTransactionState state = executeLocalTransaction(msg);
                // 4. 更新消息状态
                MessageRecord updateRecord = 
                    messageRecordMapper.selectByMsgId(msg.getKeys());
                if (state == LocalTransactionState.COMMIT_MESSAGE) {
                    updateRecord.setStatus(MessageStatus.COMMITTED);
                    updateRecord.setCommitTime(new Date());
                } else {
                    updateRecord.setStatus(MessageStatus.ROLLBACK);
                    updateRecord.setRollbackTime(new Date());
                }
                messageRecordMapper.updateStatus(updateRecord);
                return state;
            },
            record
        );
        // 5. 处理发送异常
        if (result.getLocalTransactionState() 
                == LocalTransactionState.ROLLBACK_MESSAGE) {
            // 标记失败,进行补偿
            compensateMessageRecord(record);
            throw new BusinessException("消息发送失败");
        }
    }
    private LocalTransactionState executeLocalTransaction(Message msg) {
        // 实际业务逻辑
        return LocalTransactionState.COMMIT_MESSAGE;
    }
    private void compensateMessageRecord(MessageRecord record) {
        // 补偿处理逻辑
        record.setStatus(MessageStatus.FAILED);
        record.setRetryCount(record.getRetryCount() + 1);
        messageRecordMapper.updateStatus(record);
        // 可以加入重试队列
        retryQueueProducer.send(record);
    }
}

2 死信队列处理

@Component
public class DeadLetterProcessor {
    @RocketMQMessageListener(
        topic = "%DLQ%inventory-consumer-group",
        consumerGroup = "DLQ-handler-group"
    )
    public static class DLQConsumer implements RocketMQListener<MessageExt> {
        @Override
        public void onMessage(MessageExt msg) {
            log.error("收到死信消息: {}", msg);
            // 1. 记录死信日志
            DeadLetterRecord record = DeadLetterRecord.builder()
                .msgId(msg.getMsgId())
                .topic(msg.getTopic())
                .tags(msg.getTags())
                .body(new String(msg.getBody()))
                .consumeTime(new Date())
                .build();
            deadLetterMapper.insert(record);
            // 2. 发送告警通知
            alertService.sendAlert(msg);
            // 3. 将消息重新投递到重试队列
            producer.send(new Message("RETRY-" + msg.getTopic(),
                                     msg.getBody()));
        }
    }
}

事务消息最佳实践总结

1 设计原则

原则 说明 实现方式
幂等性 多次执行结果一致 唯一索引、状态机设计
一致性 最终一致性即可 事务回查、定时补偿
可追溯 记录完整操作链路 消息ID、业务流水号
可恢复 异常自动恢复 重试队列、死信处理

2 注意事项

事务消息使用限制:

  • 适合跨进程异步同步场景
  • 无法保证强一致性
  • 需要额外实现事务回查逻辑

性能优化建议:

  • 合理设置事务消息超时时间(默认6秒)
  • 优化回查频率(默认每60秒一次)
  • 批量提交事务提升性能

监控告警配置:

  • 监控事务消息发送成功率
  • 监控本地事务提交延迟
  • 监控事务回查处理时间

生产环境注意事项

# 事务消息生产配置
transaction-producer:
  group: order-tx-producer-group
  namesrv-addr: 10.1.1.100:9876
  retry-times-wait: 3000
  send-msg-timeout: 30000
  max-retry-times: 3
# 事务超时与回查配置(RocketMQ 5.x)
broker:
  transaction-timeout: 6000          # 事务超时时间6秒
  transaction-check-interval: 60000  # 回查间隔60秒
  transaction-check-max-times: 15    # 最大回查15次

生产环境高频问题:

  1. 事务消息积压 → 合理设置消费并发、批量处理
  2. 回查失败 → 优化本地事务表设计,使用Redis辅助缓存
  3. 消息丢失 → 增加消息记录表,实现补偿机制

通过以上完整案例,您可以根据业务场景灵活实现分布式事务,保证数据最终一致性。

抱歉,评论功能暂时关闭!