RocketMQ入门教程(五):可靠消息最终一致性(事务消息)
RocketMQ的分布式事务消息是在普通消息的基础上,支持二阶段提交能力,将两阶段提交和本地事务绑定实现全局提交结果的一致性。
一:如何保证本地操作数据库和发送消息的一致性
思路一:先发消息后操作数据库
begin transaction
// 1. 发送MQ
// 2. 数据库操作
commit transaction
- MQ发送失败,当然第二步也得不到执行,保证了事务一致性。
- 可是如果是发送MQ成功了,数据库操作失败了,已经发送的MQ就收不回来了,这种情况就不能保证事务的一致性。此时奇迹出现了,RocketMQ能够实现发出去的消息再收回来。

- 生产者将消息发送至Apache RocketMQ服务端,发送这个半消息对于订阅者来说是不可见、不可消费的,必须等到步骡④才可以被消费。
- Apache RocketMO服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。
- 生产者开始执行本地事务逻辑,修改订单的状态为已支付。
- 生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rolback),服务端收到确认结果后处理逻辑如下:。
- 二次确认结果为Commit:股务端将半事务消息标记为可投递,并投递给消费者。
- 二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。
- 在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果(步骡4),或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查,回查逻辑需要自己实现,检查订单状态是否已经修改为已支付。
- 生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果,
- 生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。
事务消息生命周期:

- 初始化:半事务消息被生产者构建并完成初始化,待发送到服务端的状态。
- 事务待提交:半事务消息被发送到服务端,和普通消息不同,并不会直接被服务端持久化,而是会被单独存储到事务存储系统中,等待第二阶段本地事务返回执行结果后再提交。此时消息对下游消费者不可见。
- 消息回滚:第二阶段如果事务执行结果明确为回滚,服务端会将半事务消息回滚,该事务消息流程终止。
- 提交待消费:第二阶段如果事务执行结果明确为提交,服务端会将半事务消息重新存储到普通存储系统中,此时消息对下游消费者可见,等待被消费者获取并消费。
- 消费中:消息被消费者获取,并按照消费者本地的业务逻辑进行处理的过程。 此时服务端会等待消费者完成消费并提交消费结果,如果一定时间后没有收到消费者的响应,Apache RocketMQ会对消息进行重试处理。
- 消费提交:消费者完成消费处理,并向服务端提交消费结果,服务端标记当前消息已经被处理(包括消费成功和失败)。 ApacheRocketMQ默认支持保留所有消息,此时消息数据并不会立即被删除,只是逐辑标记已消费。消息在保存时间到期或存储空间不足被删除前,消费者仍然可以回溯消息重新消费。
- 消息删除:Apache RocketMQ按照消息保存机制滚动清理最早的消息数据,将消息从物理文件中删除。
RocketMQ支持发送一种事务消息也称为半消息,当生产者发出半消息时,此时RocketMQ会将当前消息标记为等待消费的状态,也就是说消费者是消费不到”等待消费“状态的消息的,这就是半消息的由来,生产者发了消息但是消费者收不到消息,整个过程只完成了一半。
当BrokerServer收到一个半消息时会异步通知生产者,告诉生产者MQ已经成功收到消息了(100%收到了),当生产者接受到成功发送的消息时就可以执行本地数据库的相关操作,当所有业务逻辑都成功处理完了就给MQ回一个提交状态,当MQ收到这个提交状态,就会将半消息的状态改为允许消费的状态,这样消费者就能够消费到消息。如果在执行本地数据库相关操作时出现异常就给MQ一个回滚的状态,MQ收到回滚状态就会将半消息给删除掉,这样消费者也消费不到消息,这样就保证了数据库操作和成功发送消除的一致性。
当生产者向MQ响应提交状态或者回滚状态时假如网络超时,MQ没接收到,MQ会定时调用生产者,让生产者重新响应该消息是什么状态。
二:RocketMQ发送事务消息
事务执行结果表。
CREATE TABLE `tbl_mq_tx_log` (
`id` bigint(20) NOT NULL,
`tx_no` varchar(255) DEFAULT NULL COMMENT '半事务全局唯一值',
`status` tinyint(1) DEFAULT NULL COMMENT '0: 未发送, 1: 处理成功, 2: 处理失败',
`create_time` datetime DEFAULT CURRENT_TIMESTAMP,
`update_time` datetime DEFAULT NULL,
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8
发送消息。发送消息时需要指定事务生产者组名(txProducerGroup),后面事务监听器会使用到。
@RequestMapping("/sendTx")
public void sendTx() {
String uuid = UUID.randomUUID().toString();
mqTxLogMapper.insert(new MQTxLog(uuid, 0));
Map<String, Object> headers = new HashMap<>();
headers.put("txNo", uuid);
GenericMessage message = new GenericMessage("tx msg body", headers);
TransactionSendResult transactionSendResult = rocketMQTemplate.sendMessageInTransaction("testTxProducerGroup", "tx-topic", message, uuid);
log.info("send transcation message body={},result= {}",message.getPayload(),transactionSendResult.getSendStatus());
}
本地业务逻辑处理。
@Service
public class TestService {
@Transactional(rollbackFor = Exception.class)
public void createAccount() {
System.out.println("处理业务逻辑, 插入用户");
}
}
本地事务监听器。注意rocketmq-spring-boot-starter不同版本中的@RocketMQTransactionListener参数不太一样,高版本的可能没有txProducerGroup参数,这里使用的是2.0.2。
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-spring-boot-starter</artifactId>
<version>2.0.2</version>
</dependency>
本地事务监听器。指定需要监听的事务消息组。
@Component
@RocketMQTransactionListener(txProducerGroup = "testTxProducerGroup")
public class LocalTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private MQTxLogMapper mqTxLogMapper;
@Autowired
private TestService userService;
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
System.out.println("RocketMQLocalTransactionListener#executeLocalTransaction arg=" + arg);
userService.createAccount();
// 更新事务状态为成功
MQTxLog mqTxLog = new MQTxLog();
mqTxLog.setStatus(1);
mqTxLogMapper.update(mqTxLog, Wrappers.<MQTxLog>lambdaUpdate().eq(MQTxLog::getTxNo, arg));
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
// 更新事务状态为失败
MQTxLog mqTxLog = new MQTxLog();
mqTxLog.setStatus(2);
mqTxLogMapper.update(mqTxLog, Wrappers.<MQTxLog>lambdaUpdate().eq(MQTxLog::getTxNo, arg));
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
System.out.println("消息回查");
// 在数据库中查询事务的处理结果
String txNo = msg.getHeaders().get("txNo").toString();
MQTxLog mqTxLog = mqTxLogMapper.selectOne(Wrappers.<MQTxLog>lambdaQuery().eq(MQTxLog::getTxNo, txNo));
if (mqTxLog == null) {
return RocketMQLocalTransactionState.UNKNOWN;
} else {
Integer status = mqTxLog.getStatus();
if (status == 1) {
return RocketMQLocalTransactionState.COMMIT;
} else if (status == 2) {
return RocketMQLocalTransactionState.ROLLBACK;
}
}
return RocketMQLocalTransactionState.UNKNOWN;
}
}
消费消息。
@Component
@RocketMQMessageListener(
consumerGroup = "txConsumerGroup",
topic = "tx-topic"
)
public class TxTopicConsumerListener implements RocketMQListener<String> {
@Override
public void onMessage(String messageExt) {
System.out.println(messageExt);
}
}

更多推荐
所有评论(0)