当前位置: 代码网 > it编程>编程语言>Java > SpringBoot整合RabbitMQ开发实战:工作队列、发布订阅、路由和通配符模式详解

SpringBoot整合RabbitMQ开发实战:工作队列、发布订阅、路由和通配符模式详解

2026年08月18日 Java 我要评论
对于 rabbitmq 开发,spring 也提供了一些便利。spring 和 rabbitmq 的官方文档对此均有介绍。下面来看如何基于 springboot 进行 rabbitmq 的开发。咱们只

对于 rabbitmq 开发,spring 也提供了一些便利。springrabbitmq 的官方文档对此均有介绍。

下面来看如何基于 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 均收到消息

总结

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

(0)

相关文章:

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

发表评论

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