资讯动态

一条消息的旅程:RabbitMQ 学习与实践(四)

发布时间:2026/8/28 8:09:19 来源:尧图企业网站定制
专栏RabbitMQ 进阶之路个人主页手握风云目录一、SpringBoot 整合 RabbitMQ1.1. 环境准备二、四大常用模式2.1. Work Queue 工作队列模式2.2. Publish/Subscribe 发布‑订阅2.3. Routing 路由模式2.4. Topics 通配符模式一、SpringBoot 整合 RabbitMQ1.1. 环境准备核心依赖dependencies !-- RabbitMQ整合核心依赖 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-webmvc/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp-test/artifactId scopetest/scope /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-webmvc-test/artifactId scopetest/scope /dependency /dependencies配置文件# 方式1分字段配置 spring: rabbitmq: host: 公网 IP port: 5672 username: yang password: study virtual-host: test # 方式2地址字符串简写 spring: rabbitmq: addresses: amqp://yang:study公网 IP:5672/test二、四大常用模式2.1. Work Queue 工作队列模式特点同一个队列多个消费者监听一条消息只会被其中一个消费者消费做任务分发。配置类声明队列package com.yang.rabbitmqspringboot.config; import com.yang.rabbitmqspringboot.constants.Constants; import org.springframework.amqp.core.*; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; Configuration public class RabbitMQConfig { Bean(workQueue) public Queue workQueue() { // 构建持久化队列 return QueueBuilder.durable(Constants.WORK_QUEUE).build(); } }生产者Controller 发送消息package com.yang.rabbitmqspringboot.controller; import com.yang.rabbitmqspringboot.constants.Constants; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; RestController RequestMapping(/producer) public class ProducerController { Autowired private RabbitTemplate rabbitTemplate; RequestMapping(/work) public String work() { // 向默认交换机发送 for (int i 0; i 10; i) { rabbitTemplate.convertAndSend(, Constants.WORK_QUEUE, hello spring amqp: work... i); } return 发送成功; } }消费者两个方法监听同一个队列形成竞争消费package com.yang.rabbitmqspringboot.listener; import com.rabbitmq.client.Channel; import com.yang.rabbitmqspringboot.constants.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, Channel channel) { System.out.println(Listener1 [ Constants.WORK_QUEUE ] 接收到消息 message , channel: channel); } RabbitListener(queues Constants.WORK_QUEUE) public void queueListener2(Message message, Channel channel) { System.out.println(Listener2 [ Constants.WORK_QUEUE ] 接收到消息 message , channel: channel); } }2.2. Publish/Subscribe 发布‑订阅特点fanout 交换机忽略 routingKey消息复制多份所有绑定该交换机的队列全部收到一模一样的消息。Config 配置声明 2 个队列 Fanout 交换机 两个 Binding 绑定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(fanoutQueue2) Queue queue) { return BindingBuilder.bind(queue).to(fanoutExchange); }生产者发送消息routingKey 传空字符串RequestMapping(/fanout) public String fanout() { rabbitTemplate.convertAndSend(Constants.FANOUT_EXCHANGE, , hello spring amqp: fanout...); return 发送成功; }消费者分别监听两个队列package com.yang.rabbitmqspringboot.listener; import com.yang.rabbitmqspringboot.constants.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); } }2.3. Routing 路由模式特点direct交换机routingKey 必须和 Binding 绑定的 key 完全相等队列才接收消息。适合日志分级等定向分发。Config 配置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(); } Bean(directQueueBiding1) public Binding directQueueBiding1(Qualifier(directExchange) DirectExchange directExchange, Qualifier(directQueue1) Queue queue) { return BindingBuilder.bind(queue).to(directExchange).with(orange); } Bean(directQueueBiding2) public Binding directQueueBiding2(Qualifier(directExchange) DirectExchange directExchange, Qualifier(directQueue2) Queue queue) { return BindingBuilder.bind(queue).to(directExchange).with(black); } Bean(directQueueBiding3) public Binding directQueueBiding3(Qualifier(directExchange) DirectExchange directExchange, Qualifier(directQueue2) Queue queue) { return BindingBuilder.bind(queue).to(directExchange).with(orange); }生产者routingKey 动态传入RequestMapping(/direct/{routingKey}) public String direct(PathVariable(routingKey) String routingKey) { rabbitTemplate.convertAndSend(Constants.DIRECT_EXCHANGE, routingKey, hello spring amqp: direct... routingKey); return 发送成功; }消费者监听两个队列package com.yang.rabbitmqspringboot.listener; import com.yang.rabbitmqspringboot.constants.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); } }2.4. Topics 通配符模式特点topic交换机支持通配符匹配。*匹配一个单词.分割#匹配0 个或者多个单词。Config 配置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(); } Bean(topicQueueBinding1) public Binding topicQueueBinding1(Qualifier(topicExchange) TopicExchange topicExchange, Qualifier(topicQueue1) Queue queue) { return BindingBuilder.bind(queue).to(topicExchange).with(*.orange.*); } Bean(topicQueueBinding2) public Binding topicQueueBinding2(Qualifier(topicExchange) TopicExchange topicExchange, Qualifier(topicQueue2) Queue queue){ return BindingBuilder.bind(queue).to(topicExchange).with(*.*.rabbit); } 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) { rabbitTemplate.convertAndSend(Constants.TOPIC_EXCHANGE, routingKey, hello spring amqp: topic... routingKey); return 发送成功; }消费者package com.yang.rabbitmqspringboot.listener; import com.yang.rabbitmqspringboot.constants.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); } }

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价