对于 rabbitmq 开发,spring 也提供了一些便利。spring 和 rabbitmq 的官方文档对此均有介绍。
下面来看如何基于 springboot 进行 rabbitmq 的开发。咱们只演示部分常用的工作模式:
- 工作队列模式(work queues)
- 发布订阅模式(publish / subscribe)
- 路由模式(routing)
- 通配符模式(topics)
1. 创建 springboot 项目
新建一个空的项目如下所示:

然后添加依赖:

创建好以后,把下面几个没有用的给删除掉

并且依赖已经是添加好了的

然后添加 rabbitmq 的服务配置,这里使用 yml 的格式
写法一
#配置rabbitmq的基本信息
spring:
rabbitmq:
host:
port: 5672 #默认为5672
username:
password:
virtual-host: #默认值为/或者还有下面这种写法(推荐)
#配置rabbitmq的基本信息
rabbitmq:
#amqp://username:password@ip:port/virtual-host
addresses: amqp://edison:edison@ip:5672/my_app_vhost注意:rabbitmq 的通信端口是 5672,管理后台端口是 15672。
下面就可以开始完成代码的编写了。
2. 工作队列模式
步骤:
- 1、引入依赖
- 2、编写 yml 配置,基本信息配置
- 3、编写生产者代码
- 4、编写消费者代码
- a、定义监听类,使用 @rabbitlistener 注解完成队列监听
- 5、运行观察结果
2.1 编写生产者代码
先把队列名称定义为一个常量
public static final string work_queue = "workqueue";
声明队列
package com.edison.rabbitmq.config;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.core.queue;
import org.springframework.amqp.core.queuebuilder;
import org.springframework.context.annotation.bean;
import org.springframework.context.annotation.configuration;
@configuration
public class rabbitmqconfig {
// 1. 工作模式队列
@bean("workqueue")
public queue workqueue() {
return queuebuilder.durable(constants.work_queue).build();
}
}
为了方便测试,我们通过接口来发送消息
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.rabbit.core.rabbittemplate;
import org.springframework.beans.factory.annotation.autowired;
import org.springframework.web.bind.annotation.requestmapping;
import org.springframework.web.bind.annotation.restcontroller;
@requestmapping("/producer")
@restcontroller
public class producercontroller {
@autowired
private rabbittemplate rabbittemplate; // 可以理解为rabbitmq客户端
@requestmapping("/work")
public string work() {
// 使用内置交换机, routingkey和队列名称一致
rabbittemplate.convertandsend("", constants.work_queue, "hello spring amqp: work...");
return "发送成功";
}
}
运行代码

然后在浏览器测试接口

同时从日志中也可以看到创建了一个新的连接

并且管理后台已经有队列信息了

由此可知,当我们把 springboot 程序启动以后,它其实并没有帮我们创建队列,而是在我们发送消息的时候,才会进行创建。
2.2 编写消费者代码
定义监听类
package com.edison.rabbitmq.listener;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.core.message;
import org.springframework.amqp.rabbit.annotation.rabbitlistener;
import org.springframework.stereotype.component;
@component
public class worklistener {
@rabbitlistener(queues = constants.work_queue)
public void queuelistener(message message) {
system.out.println("["+constants.work_queue+"] 接收到消息: "+message);
}
}
@rabbitlistener 是 spring 框架中用于监听 rabbitmq 队列的注解,通过使用这个注解,可以定义一个方法,以便从 rabbitmq 队列中接收消息。该注解支持多种参数类型,这些参数类型代表了从 rabbitmq 接收到的消息和相关信息。
以下是一些常用的参数类型:
- 1、string:返回消息的内容
- 2、message(
org.springframework.amqp.core.message):spring amqp 的 message 类,返回原始的消息体以及消息的属性,如消息 id,内容,队列信息等。 - 3、channel(
com.rabbitmq.client.channel):rabbitmq 的通道对象,可以用于进行更高级的操作,如手动确认消息。
消费者测试,打印消息内容

消息内容如下所示:
[workqueue] 接收到消息: (body:'hello spring amqp: work...' messageproperties [headers={}, contenttype=text/plain,
contentencoding=utf-8,
contentlength=0,
receiveddeliverymode=persistent,
priority=0,
redelivered=true,
receivedexchange=,
receivedroutingkey=workqueue,
deliverytag=1,
consumertag=amq.ctag-fxo4xtjqt74j4ydbwos0gq, consumerqueue=workqueue])
并且此时在管理后台可以看到,消息已经被消费了。
2.3 编写多个消费者代码
代码如下所示:
package com.edison.rabbitmq.listener;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.core.message;
import org.springframework.amqp.rabbit.annotation.rabbitlistener;
import org.springframework.stereotype.component;
@component
public class worklistener {
@rabbitlistener(queues = constants.work_queue)
public void queuelistener1(message message) {
system.out.println("listener 1 ["+constants.work_queue+"] 接收到消息: "+message);
}
@rabbitlistener(queues = constants.work_queue)
public void queuelistener2(message message) {
system.out.println("listener 2 ["+constants.work_queue+"] 接收到消息: "+message);
}
}
然后修改生产者代码,让其一次性发送多条消息
package com.edison.rabbitmq.controller;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.rabbit.core.rabbittemplate;
import org.springframework.beans.factory.annotation.autowired;
import org.springframework.web.bind.annotation.requestmapping;
import org.springframework.web.bind.annotation.restcontroller;
@requestmapping("/producer")
@restcontroller
public class producercontroller {
@autowired
private rabbittemplate rabbittemplate; // 可以理解为rabbitmq客户端
@requestmapping("/work")
public string work() {
for (int i = 0; i < 10; i ++)
{
// 使用内置交换机, routingkey和队列名称一致
rabbittemplate.convertandsend("", constants.work_queue, "hello spring amqp: work...");
}
return "发送成功";
}
}
然后运行结果,如下所示:

3. publish/subscribe(发布订阅模式)
在发布/订阅模型中,多了一个exchange角色。
exchange常见有三种类型,分别代表不同的路由规则:
- a) fanout:广播,将消息交给所有绑定到交换机的队列(publish/subscribe模式)
- b) direct:定向,把消息交给符合指定routing key的队列(routing模式)
- c) topic:通配符,把消息交给符合routing pattern(路由模式)的队列(topics模式)
我们先来看 fanout 广播路由规则。
3.1 编写生产者代码
和简单模式的区别是:需要创建交换机,并且绑定队列和交换机。
定义队列和交换机
// 发布订阅模式 public static final string fanout_queue1 = "fanout.queue1"; public static final string fanout_queue2 = "fanout.queue2"; public static final string fanout_exchange = "fanout.exchange";
声明队列和交换机
// 2. 发布订阅模式
// 声明2个队列, 观察是否两个队列都收到了消息
@bean("fanoutqueue1")
public queue fanoutqueue1() {
return queuebuilder.durable(constants.fanout_queue1).build();
}
@bean("fanoutqueue2")
public queue fanoutqueue2() {
return queuebuilder.durable(constants.fanout_queue2).build();
}
// 声明交换机
@bean("fanoutexchange")
public fanoutexchange fanoutexchange() {
return exchangebuilder.fanoutexchange(constants.fanout_exchange).durable(true).build();
}
绑定队列和交换机
// 声明交换机和队列的绑定
@bean("fanoutqueuebinding1")
public binding fanoutqueuebinding1(@qualifier("fanoutexchange") fanoutexchange fanoutexchange, @qualifier("fanoutqueue1") queue queue) {
return bindingbuilder.bind(queue).to(fanoutexchange);
}
@bean("fanoutqueuebinding2")
public binding fanoutqueuebinding2(@qualifier("fanoutexchange") fanoutexchange fanoutexchange, @qualifier("fanoutqueue1") queue queue) {
return bindingbuilder.bind(queue).to(fanoutexchange);
}
然后使用接口发送消息
@requestmapping("/fanout")
public string fanout() {
// routingkey为空, 表示所有队列都可以收到消息
rabbittemplate.convertandsend(constants.fanout_exchange, "", "hello spring amqp: fanout...");
return "发送成功";
}
然后启动程序,观察结果,可以看到队列已经和交换机绑定好了

然后再通过接口发送消息,可以看到此时队列中已经有了消息

3.2 编写消费者代码
交换机和队列的绑定关系及声明已经在生产方写完,所以消费者不需要再写了。
定义监听类,处理接收到的消息即可。
package com.edison.rabbitmq.listener;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.rabbit.annotation.rabbitlistener;
import org.springframework.stereotype.component;
@component
public class fanoutlistener {
@rabbitlistener(queues = constants.fanout_queue1)
public void queuelistener1(string message) {
system.out.println("队列["+constants.fanout_queue1+"] 接收到消息: "+message);
}
@rabbitlistener(queues = constants.fanout_queue2)
public void queuelistener2(string message) {
system.out.println("队列["+constants.fanout_queue2+"] 接收到消息: "+message);
}
}
3.3 运行程序
先运行项目,调用接口 http://127.0.0.1:8080/producer/fanout 发送消息
然后监听类收到消息,并打印

4. routing(路由模式)
交换机类型为 direct 时,会把消息交给符合指定 routing key 的队列。
队列和交换机的绑定,不是任意的绑定了,而是要指定一个 routingkey(路由 key)
消息的发送方在向 exchange 发送消息时,也需要指定消息的 routingkey。
exchange 也不再把消息交给每一个绑定的 key,而是根据消息的 routingkey 进行判断,只有队列的 routingkey 和消息的 routingkey 完全一致,才会接收到消息。
4.1 编写生产者代码
和发布订阅模式的区别就是:交换机类型不同,绑定队列的 routingkey 不同。
定义队列和交换机
// 路由模式 public static final string direct_queue1 = "direct.queue1"; public static final string direct_queue2 = "direct.queue2"; public static final string direct_exchange = "direct.exchange";
声明队列和交换机
// 声明队列
@bean("directqueue1")
public queue directqueue1() {
return queuebuilder.durable(constants.direct_queue1).build();
}
@bean("directqueue2")
public queue directqueue2() {
return queuebuilder.durable(constants.direct_queue2).build();
}
// 声明交换机
@bean("directexchange")
public directexchange directexchange() {
return exchangebuilder.directexchange(constants.direct_exchange).durable(true).build();
}
绑定队列和交换机
// 队列和交换机绑定
// 队列1绑定orange
@bean("directqueuebinding1")
public binding directqueuebinding1(@qualifier("directexchange") directexchange directexchange, @qualifier("directqueue1") queue queue) {
return bindingbuilder.bind(queue).to(directexchange).with("orange");
}
// 队列2绑定black,orange
@bean("directqueuebinding2")
public binding directqueuebinding2(@qualifier("directexchange") directexchange directexchange, @qualifier("directqueue2") queue queue) {
return bindingbuilder.bind(queue).to(directexchange).with("black");
}
@bean("directqueuebinding3")
public binding directqueuebinding3(@qualifier("directexchange") directexchange directexchange, @qualifier("directqueue2") queue queue) {
return bindingbuilder.bind(queue).to(directexchange).with("orange");
}
绑定关系如下图所示:

使用接口发送消息
@requestmapping("/direct/{routingkey}")
public string direct(@pathvariable("routingkey") string routingkey) { // 从路径中拿到routingkey, 需要使用pathvariable注解
// routingkey作为参数传递
rabbittemplate.convertandsend(constants.direct_exchange, routingkey, "hello spring amqp: direct, my routing key is " + routingkey);
return "发送成功";
}
然后运行程序,从管理平台上可以看到,此时队列没有任何消息

然后绑定关系如下

此时,我们通过接口 http://127.0.0.1:8080/producer/direct/orange 来发送 orange 消息

然后再通过接口 http://127.0.0.1:8080/producer/direct/black 发送 balck 消息,该消息只能由队列 2 拿到

4.2 编写消费者代码
交换机和队列的绑定关系及声明已经在生产方写完,所以消费者不需要再写了。
定义监听类,处理接收到的消息即可。
package com.edison.rabbitmq.listener;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.rabbit.annotation.rabbitlistener;
import org.springframework.stereotype.component;
@component
public class directlistener {
@rabbitlistener(queues = constants.direct_queue1)
public void queuelistener1(string message) {
system.out.println("队列["+constants.direct_queue1+"] 接收到消息: "+message);
}
@rabbitlistener(queues = constants.direct_queue2)
public void queuelistener2(string message) {
system.out.println("队列["+constants.direct_queue2+"] 接收到消息: "+message);
}
}
4.3 运行程序
运行项目
先调用接口 http://127.0.0.1:8080/producer/direct/orange 发送 routingkey 为 orange 的消息
观察后端日志,队列 1 和队列 2 收到消息

然后调用接口 http://127.0.0.1:8080/producer/direct/black 发送 routingkey 为 black 的消息
观察后端日志,队列 2 收到消息

5. 通配符模式
主题(topics)和路由(routing)模式的区别是:
- 主题模式使用的交换机类型为主题(topic)(路由模式使用的交换机类型为直连(direct))
- 主题类型的交换机在匹配规则上进行了扩展,绑定键(binding key)支持通配符匹配
5.1 编写生产者代码
和发布订阅模式的区别就是:交换机类型不同,绑定队列的路由键(routingkey)不同。
定义队列和交换机
// 通配符模式 public static final string topic_queue1 = "topic.queue1"; public static final string topic_queue2 = "topic.queue2"; public static final string topic_exchange = "topic.exchange";
声明队列和交换机
// 声明队列
@bean("topicqueue1")
public queue topicqueue1() {
return queuebuilder.durable(constants.topic_queue1).build();
}
@bean("topicqueue2")
public queue topicqueue2() {
return queuebuilder.durable(constants.topic_queue2).build();
}
// 声明交换机
@bean("topicexchange")
public topicexchange topicexchange() {
return exchangebuilder.topicexchange(constants.topic_exchange).durable(true).build();
}
绑定队列和交换机
// 队列和交换机绑定
// 队列1绑定[*.orange.*]
@bean("topicqueuebinding1")
public binding topicqueuebinding1(@qualifier("topicexchange") topicexchange topicexchange, @qualifier("topicqueue1") queue queue) {
return bindingbuilder.bind(queue).to(topicexchange).with("*.orange.*");
}
// 队列2绑定[*.*.rabbit]
@bean("topicqueuebinding2")
public binding topicqueuebinding2(@qualifier("topicexchange") topicexchange topicexchange, @qualifier("topicqueue2") queue queue) {
return bindingbuilder.bind(queue).to(topicexchange).with("*.*.rabbit");
}
// 队列2绑定[lazy.#]
@bean("topicqueuebinding3")
public binding topicqueuebinding3(@qualifier("topicexchange") topicexchange topicexchange, @qualifier("topicqueue2") queue queue) {
return bindingbuilder.bind(queue).to(topicexchange).with("lazy.#");
}
绑定关系如下图所示:

然后使用接口发送消息
@requestmapping("/topic/{routingkey}")
public string topic(@pathvariable("routingkey") string routingkey) { // 从路径中拿到routingkey, 需要使用pathvariable注解
// routingkey作为参数传递
rabbittemplate.convertandsend(constants.topic_exchange, routingkey, "hello spring amqp: topic, my routing key is " + routingkey);
return "发送成功";
}
然后运行程序,从管理平台上可以看到,此时队列没有任何消息

然后绑定关系如下

5.2 编写消费者代码
定义监听类,处理接收到的消息。
package com.edison.rabbitmq.listener;
import com.edison.rabbitmq.constant.constants;
import org.springframework.amqp.rabbit.annotation.rabbitlistener;
import org.springframework.stereotype.component;
@component
public class topiclistener {
@rabbitlistener(queues = constants.topic_queue1)
public void queuelistener1(string message) {
system.out.println("队列["+constants.topic_queue1+"] 接收到消息: "+message);
}
@rabbitlistener(queues = constants.topic_queue2)
public void queuelistener2(string message) {
system.out.println("队列["+constants.topic_queue2+"] 接收到消息: "+message);
}
}
5.3 运行程序
运行项目
调用接口 http://127.0.0.1:8080/producer/topic/green.orange.c 发送 routingkey 为 green.orange.c 的消息
观察后端日志,队列 1 收到消息

调用接口 http://127.0.0.1:8080/producer/topic/green.orange.rabbit 发送 routingkey 为 green.orange.rabbit 的消息
观察后端日志,队列 1 和队列 2 均收到消息

调用接口 http://127.0.0.1:8080/producer/topic/lazy.orange.rabbit 发送 routingkey 为 lazy.orange.rabbit 的消息
观察后端日志,队列 1 和队列 2 均收到消息

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