本文目录导读:

我将为您详细介绍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次
生产环境高频问题:
- 事务消息积压 → 合理设置消费并发、批量处理
- 回查失败 → 优化本地事务表设计,使用Redis辅助缓存
- 消息丢失 → 增加消息记录表,实现补偿机制
通过以上完整案例,您可以根据业务场景灵活实现分布式事务,保证数据最终一致性。