当前位置: 代码网 > it编程>编程语言>Java > SpringBoot集成RocketMQ常见问题有哪些?版本冲突和监听器配置详解

SpringBoot集成RocketMQ常见问题有哪些?版本冲突和监听器配置详解

2026年09月06日 Java 我要评论
一、原生rocketmq配置复杂,需要手动配置product、consumer、监听器、序列化、消息过滤、事物消息;非常灵活,所有rocketmq底层能力都能用上。使用建议: 对 rocketmq 特

一、原生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属性,实现了消息溯源、消息压缩等⼀系列功能。

总结

以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。

(0)

相关文章:

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com