一、原生rocketmq
配置复杂,需要手动配置product、consumer、监听器、序列化、消息过滤、事物消息;
非常灵活,所有rocketmq底层能力都能用上。
使用建议: 对 rocketmq 特性依赖较深(如事务消息、消息延迟等级、过滤等)。 项目中对性能和精细化控制要求高。
1.添加依赖:
<dependency>
<groupid>org.apache.rocketmq</groupid>
<artifactid>rocketmq-client</artifactid>
<version>5.3.0</version>
</dependency>2.生产者:
public class producer {
public static void main(string[] args) throws mqclientexception, interruptedexception {
//1.初始化一个消息生产者,指定组名
defaultmqproducer producer = new defaultmqproducer("demoproducer");
// 2.指定nameserver地址
producer.setnamesrvaddr("192.168.65.112:9876");
// 3.启动消息生产者服务
producer.start();
for (int i = 0; i < 2; i++) {
try {
// 4.创建消息。消息由topic,tag和body三个属性组成,其中body就是消息内容
message msg = new message("topictest","taga",("hello rocketmq " +i).getbytes(remotinghelper.default_charset));
//5.发送消息,获取发送结果
sendresult sendresult = producer.send(msg);
system.out.printf("%s%n", sendresult);
} catch (exception e) {
e.printstacktrace();
thread.sleep(1000);
}
}
//6.消息发送完后,停止消息生产者服务。
producer.shutdown();
}
}3.消费者:
public class consumer {
public static void main(string[] args) throws interruptedexception, mqclientexception {
//1.构建一个消息消费者指定消费者组
defaultmqpushconsumer consumer = new defaultmqpushconsumer("please_rename_unique_group_name_4");
//2.指定nameserver地址
consumer.setnamesrvaddr("192.168.65.112:9876");
consumer.setconsumefromwhere(consumefromwhere.consume_from_last_offset);
// 3.订阅一个感兴趣的话题,这个话题需要与消息的topic一致,不实用tag过滤
consumer.subscribe("topictest", "*");
// 4.注册一个消息回调函数,消费到消息后就会触发回调。
consumer.registermessagelistener(new messagelistenerconcurrently() {
@override
public consumeconcurrentlystatus consumemessage(list<messageext> msgs,consumeconcurrentlycontext context) {
msgs.foreach(messageext -> {
try {
system.out.println("收到消息:"+new string(messageext.getbody(), remotinghelper.default_charset));
} catch (unsupportedencodingexception e) {}
});
//返回消费状态
return consumeconcurrentlystatus.consume_success;
}
});
//5.启动消费者服务
consumer.start();
system.out.print("consumer started");
}
}二、spring-boot-starter-rocketmq
使用 @rocketmqmessagelistener 和 rocketmqtemplate 快速开发,封装好 producer 和 consumer 配置简单,springboot自动装配,不够灵活
1.添加依赖(使用springboot集成时,版本特别关键,稍微有偏差就会报错)
<dependencies> <dependency> <groupid>org.apache.rocketmq</groupid> <artifactid>rocketmq-spring-boot-starter</artifactid> <version>2.3.1</version> <exclusions> <exclusion> <groupid>org.apache.rocketmq</groupid> <artifactid>rocketmq-client</artifactid> </exclusion> </exclusions> </dependency> <dependency> <groupid>org.apache.rocketmq</groupid> <artifactid>rocketmq-client</artifactid> <version>5.3.0</version> </dependency> <dependency> <groupid>org.springframework.boot</groupid> <artifactid>spring-boot-starter-web</artifactid> <version>3.0.4</version> </dependency> <dependency> <groupid>org.springframework.boot</groupid> <artifactid>spring-boot-starter-test</artifactid> <version>3.0.4</version> </dependency> <dependency> <groupid>junit</groupid> <artifactid>junit</artifactid> <version>4.13.2</version> <scope>test</scope> </dependency> </dependencies>
2.配置文件
rocketmq.name-server=192.168.65.112:9876 rocketmq.producer.group=springbootgroup #如果这⾥不配,那就需要在消费者的注解中配。 #rocketmq.consumer.topic= rocketmq.consumer.group=testgroup server.port=9000
3.生产者
rocketmqtemplate 不光可以发送消息还可以拉消息
@resource
private rocketmqtemplate rocketmqtemplate;
public void sendmessage(string topic,string msg){
this.rocketmqtemplate.convertandsend(topic,msg);
}
}4.消费者
消费者的声明也很简单。所有属性通过@rocketmqmessagelistener注解声明
@component
@rocketmqmessagelistener(consumergroup = "myconsumergroup", topic = "testtopic",consumemode=
consumemode.concurrently,messagemodel= messagemodel.broadcasting)
public class springconsumer implements rocketmqlistener<string> {
@override
public void onmessage(string message) {
system.out.println("received message : "+ message);
}
}springboot框架中对消息的封装与原⽣api的消息封装是不⼀样的
在springboot封装的rocketmq中,默认的rocketmqtemplate只能处理一开始初始化的生产者组,特别不方便,可以通过@extrocketmqtemplateconfiguration()自己在额外配置需要的。
原生的代码量多,但是更加令快,spinrgboot封装的则适合场景简单快速开发,具体以自己业务场景为主。

producer发送的message对象是没有msgid属性的。broker端接收到producer发过来的消息后,会给每条消息单独分配⼀个唯⼀的msgid。这个msgid可以作为消息的唯⼀主键来使⽤。
但是需要注意,对于客户端来说,毕竟是不知道这个msgid是如何产⽣的。实际上,在rocketmq内部,也会针对批量消息、事务消息等特殊的消息机制,有特殊的msgid分配机制。因此,在复杂业务场景下,不建议使⽤msgid来作为消息的唯⼀索引,⽽建议采⽤下⾯的key属性,⾃⾏指定业务层⾯上的唯⼀索引。

针对key这⼀个属性,建议在业务中可以添加⼀些带有业务唯⼀性的数据,作为messageid的补充。
rocketmq基于keys属性,实现了消息溯源、消息压缩等⼀系列功能。
总结
以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。
发表评论