3、RocketMQ的顺序消息、批量消息和事务消息
前面初步体验了一下RocketMQ的安装部署和基本使用,接下来简单介绍一下RocketMQ的核心特性,包括RocketMQ的有序消息、批量消息以及事务消息
1、有序消息
有序消息又叫顺序消息,是指消费者消费消息的顺序和产生消息的顺序相同。因为在有些业务场景下是要求保证消息的有序消费的,典型的电商场景下,比如订单的生成、支付、发货等这些消息同一个订单必须按照一定的先后顺序处理才行。
消息的有序包括生产者消息发送的有序、broker消息存储的有序以及消费者消费的有序,这三个环节缺一不可,RocketMQ是支持有序消息的,在RocketMQ中,顺序消息又分为全局有序消息和分区(queue)有序消息,创建topic的时候可以设定FIFO类型的topic

1.1、全局顺序消息
一个topic内所有的消息都发布到同一个队列(创建topic的时候会创建多个queue,全局顺序消息只用一个)上,按照先进先出的顺序进行发布和消费

比如订单相关的业务统一用一个topic,而且所有的消息都发送到topic的多个队列中的某一个队列,这样所有的消息就都可以按照FIFO原则进行顺序的发布和消费。但是这种方案只能针对一些性能要求不是很高的场景。
1.2、分区顺序消息
对于指定同一个Topic,所有的消息通过sharding-key进行区块(queue)分区。同一个queue里边的消息是可以保证先进先出的顺序性要求的。sharding-key是顺序消息用来分区的关键字段,和普通消息的key不同,这点类似于数据库的分库分表的路由分片键。具体的路由规则可以选择或者自定义。

rocketmq可以根据sharding key去决定消息发送到哪个queue。在一定程度上可以提升mq的消息处理性能。
1.3、全局顺序消息和分区顺序消息对比
1.3.1、消息类型对比

**总结:**顺序消息会牺牲性能,而且不支持事务消息和定时消息
1.3.2、发送方式对比

**总结:**顺序消息不支持异步发送和单向发送
1.4、如何保证消息的顺序性
前面提到,要保证消息的顺序性,需要从三个环节去考虑:
① 消息发送的时候保证消息顺序发送
② 消息在broker上存储时需要保持和消息生产的时候一样的顺序
③ 消息被消费的时候需要保持和存储的时候一样的顺序

上图演示的是全局顺序消息不管是什么类型的消息都落在了一个队列queue上,而分区顺序消息会根据sharding key(消息类型)落在不同的队列queue上
1.5、RocketMQ如果实现顺序消息(以订单为例)
用户可以使用订单ID作为sharding key,这样相同orderId的消息就都会路由到同一个队列,这样就保证了消息的顺序性,如下图所示:

RocketeMQ消费端有两种类型:MQPullConsumer和MQPushConsumer,其本质上都是通过pull长轮询机制去实现的,push是对pull的一种api封装
MQPullConsumer是由用户线程去控制的,主动从服务端获取消息,每次获取到的是一个MessageQueue中的消息。
PullResult中的List<MessageExt> msgFoundList自然和存储顺序一致,用户需要在拿到这批消息之后自己保证消费的顺序性。
MQPushConsumer是由用户注册的MessageListener来消费消息的,用户需要在客户端中保证调用MessageListener时消息的顺序性。
1.6、顺序消息demo
1.6.1、引入依赖
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>5.5.1</version>
</dependency>
1.6.2、自定义消息类
package com.study.rocketmq.order;
import java.io.Serializable;
public class CustomMessage implements Serializable {
private Integer msgType;
private String userId;
private String desc;
public CustomMessage(String desc, Integer msgType, String userId) {
this.desc = desc;
this.msgType = msgType;
this.userId = userId;
}
// 此处省去了getter、setter方法
@Override
public String toString() {
return "CustomMessage{" +
"desc='" + desc + '\'' +
", msgType=" + msgType +
", userId='" + userId + '\'' +
'}';
}
}
1.6.3、Producer
package com.study.rocketmq.order;
import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.MessageQueueSelector;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.apache.rocketmq.remoting.exception.RemotingException;
import java.io.UnsupportedEncodingException;
import java.util.ArrayList;
import java.util.List;
public class Producer {
// 定义NameServer地址,多个地址用分号隔开
public static final String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
public static void main(String[] args) throws MQClientException, UnsupportedEncodingException, MQBrokerException, RemotingException, InterruptedException {
// 1:创建生产者对象,并指定组名
DefaultMQProducer producer = new DefaultMQProducer("GROUP_TEST");
// 2:指定NameServer地址
producer.setNamesrvAddr(NAME_SERVER_ADDR);
// 3:启动生产者
producer.start();
// 设置异步发送失败重试次数,默认为2
producer.setRetryTimesWhenSendAsyncFailed(0);
// 4:定义消息队列选择器
MessageQueueSelector messageQueueSelector = new MessageQueueSelector() {
/**
* 消息队列选择器,保证同一条业务数据的消息在同一个队列
* @param mqs topic中所有队列的集合
* @param msg 发送的消息
* @param arg 此参数是本示例中producer.send的第三个参数
* @return
*/
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Integer id = (Integer) arg;
// id == 1001
int index = id % mqs.size();
// 分区顺序:同一个模值的消息在同一个队列中
return mqs.get(index);
// 全局顺序:所有的消息都在同一个队列中
// return mqs.get(0);
}
};
String[] tags = new String[]{"TagA", "TagB", "TagC"};
List<CustomMessage> bizDatas = getBizData();
// 5:循环发送消息
for (int i = 0; i < bizDatas.size(); i++) {
CustomMessage bizData = bizDatas.get(i);
// keys:业务数据的ID,比如用户ID、订单编号等等
Message msg = new Message("TopicTest", tags[i % tags.length], "" + bizData.getMsgType(), bizData.toString().getBytes(RemotingHelper.DEFAULT_CHARSET));
// 发送有序消息
SendResult sendResult = producer.send(msg, messageQueueSelector, bizData.getMsgType());
System.out.printf("%s, \n\nbody:%s%n", sendResult, bizData);
}
// 6:关闭生产者
producer.shutdown();
}
public static List<CustomMessage> getBizData() {
List<CustomMessage> orders = new ArrayList<CustomMessage>();
orders.add(new CustomMessage("下单:男款破洞鞋39码蓝色", 1, "张三"));
orders.add(new CustomMessage("下单:碎花T恤XL码", 1, "李四"));
orders.add(new CustomMessage("下单:印度神油壹号330ml", 1, "王麻子"));
orders.add(new CustomMessage("下单:304不锈钢筷子25cm(圆+圆)", 1, "老赵"));
orders.add(new CustomMessage("下单:晨光文具标准高考套装", 1, "钱六孙"));
orders.add(new CustomMessage("下单:3岁以下儿童夏季洒水枪蓝色背包款", 1, "周大炮"));
orders.add(new CustomMessage("下单:冬季居家棉拖鞋女款37码靛蓝色", 1, "吴老邪"));
orders.add(new CustomMessage("下单:招财猫储钱罐大号", 1, "郑大发"));
orders.add(new CustomMessage("下单:苹果18ProMax天青色512g", 1, "冯锡范"));
return orders;
}
}
1.6.4、Consumer
package com.study.rocketmq.order;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.io.UnsupportedEncodingException;
import java.util.List;
public class Consumer {
public static final String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
public static void main(String[] args) throws Exception {
// 1. 创建消费者(Push)对象
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
consumer.setNamesrvAddr(NAME_SERVER_ADDR);
/**
* 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费<br>
* 如果非第一次启动,那么按照上次消费的位置继续消费
*/
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
// 3. 订阅对应的主题和Tag
consumer.subscribe("TopicTest", "TagA || TagB || TagC");
// 4. 注册消息接收到Broker消息后的处理接口
consumer.registerMessageListener(new MessageListenerOrderly() {
// 按顺序消费
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
context.setAutoCommit(true);
doBiz(msgs.get(0));
return ConsumeOrderlyStatus.SUCCESS;
}
});
// 5. 启动消费者(必须在注册完消息监听器后启动,否则会报错)
consumer.start();
System.out.println("已启动消费者");
}
/**
* 模拟处理业务
*
* @param message
*/
public static void doBiz(Message message) {
try {
System.out.printf("线程:%-25s 接收到新消息 %s --- %s %n", Thread.currentThread().getName(), message.getTags(), new String(message.getBody(), RemotingHelper.DEFAULT_CHARSET));
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
}
}
1.6.5、运行
先启动消费者进行消息订阅,然后启动生产者发送消息,从控制台可以看出消息已经正常发送成功了

再看消费者端控制台输出:

可以看到消费者端打印出来的消息消费顺序和生产者发送消息的顺序是完全一致的。
不难发现,我们消费者端代码中在注册监听的时候用的是MessageListenerOrderly而不是前面用过的MessageListenerConcurrently,这是一个按顺序消费的消息监听器,RockeetMQ在底层帮我们实现了顺序消费问题,前提是只针对Topic的一个queue,而前面生产者在发送消息的时候指定了一个MessageQueueSelector,通过这个选择器的select我们按照一定的规则指定返回了一个MessageQueue,这就是前面的sharding key对应的消息队列选择器,我们可以根据我们的实际要求来实现。
1.7、顺序消息的不足
通过前面的分析不难看出,顺序消息典型的不足之处是性能没有普通消息高,特别是全局顺序消息,Topic下的所有消息都发往一个queue,操作的也都是这同一个queue,显然不如多个queue并发处理的效率高
其次,顺序消息无法利用集权的failover特性,因为消息的sharding key和sharding规则一旦确定,这个消息被分配到哪个messagequeue就确定了,无法更换MessageQueue重试,而且因为路由策略设计问题,有可能MessageQueue之间出现严重的数据倾斜
消费的并行读依赖于queue的数量,也就是一个queue只能有一个消费者来读,如果有多个来读的话就无法保证消息消费的顺序了
消费失败时也无法跳过,因为一旦没有收到明确消费成功而跳过后前面的消息又被消费了,这样一来可能就会出现后面的消息比前面的消息先消费成功
2、RocketMQ发布/订阅的基本概念
几乎所有的MQ都是发布/订阅的实现,所谓的发布订阅,在设计模式上又叫观察者模式,它定义了对象之间的一种一对多的依赖关系,当一个对象的状态发生改变的时候,所有依赖它的对象都会得到通知。就像RocketMQ消费者会订阅一个Topic,并指定匹配的标签,如:
consumer.subscribe("TopicTest", "TagA || TagB || TagC");
RocketMQ的发布订阅分为两种模式:
- Push模式(MQPushConsumer):broker主动向消费者发送消息
- Pull模式(MQPullConsumer):消费者在自己需要的时候主动向broker拉取消息
2.1、Push模式实现原理

消费者流程如上图所示:
①、消费者主题订阅注册监听
②、消费者通过pull方式从broker获取消息
③、如果拉取到有新的消息,则进行消息消费,消费完成后继续从broker拉取消息
④、如果此时broker没有新的消息,broker进行阻塞,直到有新的消息发送过来,或者消费者连接超时
⑤、消费者从broker拉取消息消费完成或拉取消息超时,继续重新拉取,依此循环
2.2、Pull模式实现原理
在pull方式中,取消息的过程是需要用户自己实现的,简单流程:
①、通过即将消费的Topic获取到绑定到这个主题下的所有MessageQueue
②、遍历这些MessageQueue集合,从每个MessageQueue中批量获取消息
③、一次取完后,记录下这个队列下一次要取的位置offset,下次从这个位置继续取消息,直到这个队列所有消息都被取完,然后继续下一个MessageQueue
3、定时消息
定时消息是指消息发送到broker后,不能立刻被consumer消费,要等到特定的时间点或者要等待特定的时间之后才能正常被消费。
在rabbitmq中,我们可以通过死信队列来实现订单的延时关闭,也可以通过一个delay插件来实现延时消息,那么在rocketmq中,如果要实现定时消息该怎么做呢?
如果要支持任意时间精度的定时,broker层面需要做消息的排序,会带来非常大的开销。RocketMQ虽然支持定时消息,但是不支持任意精度,而实内部定义了18个level级别,分别对应18个时间

3.1、定时消息发送逻辑
定时消息的发送逻辑如下图所示:

①、首先获取消息队列,如果获取不到队列,则下次定时任务再执行
②、如果获取到了消息队列,就从队列中获取消息
③、如果也没有获取到消息,则下次定时任务再执行
④、如果获取到了新的消息,则计算消息的时间(根据lever),判断是否达到了发送消息的时候指定的延时时间
⑤、如果消息达到了正常延时的时间,则进行发送,否则安排后续定时任务n 100ms执行(定时任务衰减重试)
⑥、消息发送成功,则继续遍历队列,如果消息发送失败,则安排下次定时任务10000ms后执行
3.2、定时消息代码示例
Producer.java
package com.study.rocketmq.scheduler;
import org.apache.commons.lang3.RandomUtils;
import org.apache.rocketmq.client.exception.MQBrokerException;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import org.apache.rocketmq.remoting.exception.RemotingException;
import java.io.UnsupportedEncodingException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.concurrent.TimeUnit;
public class Producer {
// 定义NameServer地址,多个地址用分号隔开
public static final String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
public static void main(String[] args) throws MQClientException, InterruptedException, RemotingException, MQBrokerException, UnsupportedEncodingException {
// 1. 创建生产者对象
DefaultMQProducer producer = new DefaultMQProducer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
producer.setNamesrvAddr(NAME_SERVER_ADDR);
// 3. 启动生产者
producer.start();
for (int i = 0; i < 10; i++) {
String content = "Hello scheduled message " + new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SS").format(new Date());
Message message = new Message("TopicTest", content.getBytes(RemotingHelper.DEFAULT_CHARSET));
// 4. 设置延时等级,此消息将在10秒后传递给消费者
// 可以在broker服务器端自行配置messageDelayLevel=1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
message.setDelayTimeLevel(3);
// 5. 发送消息
SendResult result = producer.send(message);
System.out.printf("发送结果:%s%n", result);
TimeUnit.MILLISECONDS.sleep(RandomUtils.nextInt(300, 800));
}
// 6. 停止生产者
producer.shutdown();
}
}
Consumer.java
package com.study.rocketmq.scheduler;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.*;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.io.UnsupportedEncodingException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.List;
public class Consumer {
public static final String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
public static void main(String[] args) throws MQClientException {
// 1. 创建消费者(Push)对象
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
consumer.setNamesrvAddr(NAME_SERVER_ADDR);
// 3. 订阅对应的主题和Tag
consumer.subscribe("TopicTest", "*");
// 4. 注册消息接收到Broker消息后的处理接口
consumer.registerMessageListener(new MessageListenerConcurrently() {
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
MessageExt messageExt = list.get(0);
try {
System.out.printf("%-25s 接收到新消息 --- %s %n", new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SS").format(new Date()), new String(messageExt.getBody(), RemotingHelper.DEFAULT_CHARSET));
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
// 5. 启动消费者(必须在注册完消息监听器后启动,否则会报错)
consumer.start();
System.out.println("已启动消费者");
}
}
执行结果如图所示:
生产者成功发送10条消息:

消费者延迟大约10s后收到消息,因为我们在发送消息的时候设置的延迟等级是3,对应也就是10s:

4、批量消息
为了提升RocketMQ的吞吐能力,RocketMQ会在消息数量特别大的时候可以进行批量处理。比如消费者在消费消息时候每次只是从broker上获取一条消息,这网络连接的建立和关闭开销比数据传输的成本都高,如果消息数量足够多的话为什么不一次多拿几条消息呢?
4.1、消费端如何开启批量消费呢
要开启批量获取消息,RocketMQ提供了一个配置:
consumer.setConsumeMessageBatchMaxSize(10); // 默认32,在broker.properties的maxTransferCountOnMessageInMemory配置参数决定,超过这个值无效
consumer.setPullBatchSize(10)
这个参数默认是1,也就是单次只获取一条消息。我们可以适当调整一下这个参数的大小,这样在消费者监听器里边List<MessageExt> list一次可能得到的就是一批消息了:
consumer.registerMessageListener(new MessageListenerConcurrently() {
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> list, ConsumeConcurrentlyContext consumeConcurrentlyContext) {
try {
// 设置消息批处理数量后,list中才会有多条,否则每次只会有一条
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
4.2、批量消息带来的弊端
虽然批量消息在一定程度上会带来性能上的提升,但同时也会引起一些其它方便的弊端,比如:
① 同一批次的消息应该具有相同的主题、相同的消息配置
② 批量消息不支持延时发送
③ 单个批次消息建议大小不要超过1mb,一次拉取的消息太大,网络传输会比较慢
5、事务消息
RocketMQ的事务消息,是指Producer端发送消息事件和本地事务事件,同时成功或失败。
5.1、事务消息原理
RocketMQ是支持事务消息的,RocketMQ的事务消息原理如图所示:

过程如下:
1、生产者将消息发送至Apache RocketMQ服务端。
2、Apache RocketMQ服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。
3、生产者开始执行本地事务逻辑。
4、生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rollback),服务端收到确认结果后处理逻辑如下:
二次确认结果为Commit:服务端将半事务消息标记为可投递,并投递给消费者。
二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。
5、在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果,或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查。 说明 服务端回查的间隔时间和最大回查次数,请参见参数限制:https://rocketmq.apache.org/zh/docs/introduction/03limits。
6、生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果。
7、生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。
5.2、如何使用事务消息
从上面RoecktMQ的事务消息模型图可以看出它的事务消息主要体现在生产者端,消费端的消费逻辑和普通消息没有什么区别
5.2.1、事务消息的三种状态
- TransactionStatus.CommitTransaction:提交事务。消费者可以正常消费到这条消息
- TransactionStatus.RollbackTransaction:回滚事务,消息将会被删除或消费者不在允许消费这条消息
- TransactionStatus.Unknown:未知状态,MQ需要检查这条消息对应的本地事务来确认消息的状态
5.2.2、生产者
之前在发送消息的时候,我们使用的都是DefaultMQProducer,如果要使用事务消息,需要用到TransactionMQProducer,并且指定一个事务监听器
producer.setTransactionListener(transactionListenerImpl);
5.2.3、事务监听器
从事务消息的原理图可以看到有两个方法肯定需要我们手动去实现:执行本地事务和检查本地事务状态,这也是前面设置的事务监听器实现类需要实现的两个方法:
private final static TransactionListener transactionListenerImpl = new TransactionListener() {
/**
* 在发送消息成功时执行本地事务
* @param msg message
* @param arg producer.sendMessageInTransaction的第二个参数
* @return 返回事务状态
* LocalTransactionState.COMMIT_MESSAGE:提交事务,提交后broker才允许消费者使用
* LocalTransactionState.RollbackTransaction:回滚事务,回滚后消息将被删除,并且不允许别消费
* LocalTransactionState.Unknown:中间状态,表示MQ需要核对,以确定状态
*/
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// TODO 开启本地事务(实际就是我们的jdbc操作)
// 根据业务执行结果决定是回滚还是提交
/*
return LocalTransactionState.ROLLBACK_MESSAGE;
return LocalTransactionState.UNKNOW;
*/
return LocalTransactionState.COMMIT_MESSAGE;
}
/**
* Broker端对未确定状态的消息发起回查,将消息发送到对应的Producer端(同一个Group的Producer),
* 由Producer根据消息来检查本地事务的状态,进而执行Commit或者Rollback
* @param msg
* @return 返回事务状态
*/
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 根据业务,正确处理: 订单场景,只要数据库有了这条记录,消息应该被commit
String transactionId = msg.getTransactionId();
String key = msg.getKeys();
System.out.printf("回查事务状态 key:%-5s msgId:%-10s transactionId:%-10s %n", key, msg.getMsgId(), transactionId);
// 检查本地事务状态,决定是提交还是回滚
// return LocalTransactionState.COMMIT_MESSAGE;
return LocalTransactionState.ROLLBACK_MESSAGE;
}
};
5.2.4、消息发送
普通消息的发送用的是producer.send,事务消息的发送稍有不同:
SendResult result = producer.sendMessageInTransaction(message, arg);
第一个参数很好理解,就是消息本身,第二个参数对应事务监听器实现类executeLocalTransaction方法的第二个参数
5.2.5、生产者完整demo
package com.study.rocketmq.transaction;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.common.RemotingHelper;
import java.io.IOException;
public class Producer {
public static final String NAME_SERVER_ADDR = "192.168.0.61:9876;192.168.0.62:9876;192.168.0.63:9876";
/**
* 事务消息监听实现
*/
private final static TransactionListener transactionListenerImpl = new TransactionListener() {
/**
* 在发送消息成功时执行本地事务
* @param msg message
* @param arg producer.sendMessageInTransaction的第二个参数
* @return 返回事务状态
* LocalTransactionState.COMMIT_MESSAGE:提交事务,提交后broker才允许消费者使用
* LocalTransactionState.RollbackTransaction:回滚事务,回滚后消息将被删除,并且不允许别消费
* LocalTransactionState.Unknown:中间状态,表示MQ需要核对,以确定状态
*/
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// TODO 开启本地事务(实际就是我们的jdbc操作)
// TODO 执行业务代码(插入订单数据库表)
// int i = orderDatabaseService.insert(....)
// TODO 提交或回滚本地事务(如果用spring事务注解,这些都不需要我们手工去操作)
// 模拟一个处理结果
int index = 8;
/**
* 模拟返回事务状态
*/
switch (index) {
case 3:
System.out.printf("本地事务回滚,回滚消息,id:%s%n", msg.getKeys());
return LocalTransactionState.ROLLBACK_MESSAGE;
case 5:
case 8:
return LocalTransactionState.UNKNOW;
default:
System.out.println("事务提交,消息正常处理");
return LocalTransactionState.COMMIT_MESSAGE;
}
}
/**
* Broker端对未确定状态的消息发起回查,将消息发送到对应的Producer端(同一个Group的Producer),
* 由Producer根据消息来检查本地事务的状态,进而执行Commit或者Rollback
* @param msg
* @return 返回事务状态
*/
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 根据业务,正确处理: 订单场景,只要数据库有了这条记录,消息应该被commit
String transactionId = msg.getTransactionId();
String key = msg.getKeys();
System.out.printf("回查事务状态 key:%-5s msgId:%-10s transactionId:%-10s %n", key, msg.getMsgId(), transactionId);
if ("id_5".equals(key)) { // 刚刚测试的10条消息, 把id_5这条消息提交,其他的全部回滚。
System.out.printf("回查到本地事务已提交,提交消息,id:%s%n", msg.getKeys());
return LocalTransactionState.COMMIT_MESSAGE;
} else {
System.out.printf("未查到本地事务状态,回滚消息,id:%s%n", msg.getKeys());
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
};
public static void main(String[] args) throws MQClientException, IOException {
// 1. 创建事务生产者对象
// 和普通消息生产者有所区别,这里使用的是TransactionMQProducer
TransactionMQProducer producer = new TransactionMQProducer("GROUP_TEST");
// 2. 设置NameServer的地址,如果设置了环境变量NAMESRV_ADDR,可以省略此步
producer.setNamesrvAddr(NAME_SERVER_ADDR);
// 3. 设置事务监听器
producer.setTransactionListener(transactionListenerImpl);
// 4. 启动生产者
producer.start();
for (int i = 0; i < 10; i++) {
String content = "Hello transaction message " + i;
Message message = new Message("TopicTest", "TagA", "id_" + i, content.getBytes(RemotingHelper.DEFAULT_CHARSET));
// 5. 发送消息(发送一条新订单生成的通知)
SendResult result = producer.sendMessageInTransaction(message, i);
System.out.printf("发送结果:%s%n", result);
}
System.in.read();
// 6. 停止生产者
producer.shutdown();
}
}
注意:
RocketMQ 事务消息保证本地主分支事务和下游消息发送事务的一致性,但不保证消息消费结果和上游事务的一致性。因此需要下游业务分支自行保证消息正确处理,且事务消息为最终一致性,即在消息提交到下游消费端处理完成之前,下游分支和上游事务之间的状态会不一致。因此,事务消息仅适合接受异步执行的事务场景
5.3、事务消息的局限性
事务消息在解决了一致性问题的同时,也会带来一定的不足,主要体现在下面几个方面:
① 事务消息不支持定时和批量
② 一个消息多次检查,会增加队列的堆积概率,rocketmq为了缓解这个问题,限制了同一个消息检查的次数,这个值在broker.properties文件中可以配置:transactionMaxCheck,默认是15。与此同时还可以配置指定特定时间开始检查事务状态,这个由参数"transactionTimeout"或者"CHECK_IMMUNITY_TIME_IN_SECOND"

