在上篇文章 springboot 整合 rabbitmq 中,我们实现了springboot 和rabbitmq的简单整合,这篇文章主要是对上篇文章功能的增强,主要完成如下功能。
需求:
生产者在启动的时候,自动创建好队列、绑定、交换器并设置好 死信交换器、备份交换器(alternate-exchange)。生产者发送消息后,生产者这边需要对发送的消息进行确认,确认RabbitMQ接收到了消息。为了测试未被路由的消息和死信消息,发送方,发送11条正常的,可以被路由到消息队列中的消息,发送一条不可路由到消息队列中的消息,使之进入 alternate-exchange 交换器中。接收方在接收到消息后,随机拒绝一些消息,使之进入 x-dead-letter-exchange 中。
实现如下功能:
1、使用@Bean方式自动实现队列、交换器、绑定的创建。
2、使用@RabbitListener实现队列消息的监听。
3、实现生产者消息确认。
4、实现死信交换器(过期的消息、basic.nack或basic.reject且requeue参数为false或队列满的消息将进入此交换器)。
5、实现备份交换器(alternate-exchange),未被正确路由的消息将会经过此交换器。
部分功能实现要点:
1、生产者消息确认
|- spring.rabbitmq.template.mandatory = true 设置成true
|- spring.rabbitmq.publisher-confirms = true 设置成true
|- 编写一个 java 类,实现 RabbitTemplate.ConfirmCallback 接口,在这个里面我们可以确认消息是否到达了RabbitMQ服务器。
2、死信交换器的实现
|- 申明队列的时候设置 x-dead-letter-exchange 参数
3、处理未被路由的消息
|- 申明交换器的时候设置 alternate-exchange 参数
实现步骤如下:
1、jar包的引入,生产者和消费者都一样
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> </dependencies>
2、生产者 - 配置文件
server: port: 9088 spring: rabbitmq: host: 140.143.237.224 port: 5672 username: root password: root virtual-host: / connection-timeout: 10000 template: mandatory: true publisher-confirms: true
注意:此处需要将 mandatory和publisher-confirms参数设置成true
3、生产者 - 生产者消息确认编写
@Slf4j public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback { @Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { log.info("(start)生产者消息确认========================="); log.info("correlationData:[{}]", correlationData); log.info("ack:[{}]", ack); log.info("cause:[{}]", cause); if (!ack) { log.info("消息可能未到达rabbitmq服务器"); } log.info("(end)生产者消息确认========================="); } }
注意:此处需要实现 RabbitTemplate.ConfirmCallback接口,实现消息的确认
4、生产者 - 生产者配置
@Configuration public class RabbitmqConfiguration { @Autowired private RabbitTemplate rabbitTemplate; @PostConstruct public void initRabbitTemplate() { // 设置生产者消息确认 rabbitTemplate.setConfirmCallback(new RabbitConfirmCallback()); } /** * 申明队列 * * @return */ @Bean public Queue queue() { Map<String, Object> arguments = new HashMap<>(4); // 申明死信交换器 arguments.put("x-dead-letter-exchange", "exchange-dlx"); return new Queue("queue-rabbit-springboot-advance", true, false, false, arguments); } /** * 没有路由到的消息将进入此队列 * * @return */ @Bean public Queue unRouteQueue() { return new Queue("queue-unroute"); } /** * 死信队列 * * @return */ @Bean public Queue dlxQueue() { return new Queue("dlx-queue"); } /** * 申明交换器 * * @return */ @Bean public Exchange exchange() { Map<String, Object> arguments = new HashMap<>(4); // 当发往exchange-rabbit-springboot-advance的消息,routingKey和bindingKey没有匹配上时,将会由exchange-unroute交换器进行处理 arguments.put("alternate-exchange", "exchange-unroute"); return new DirectExchange("exchange-rabbit-springboot-advance", true, false, arguments); } @Bean public FanoutExchange unRouteExchange() { // 此处的交换器的名字要和 exchange() 方法中 alternate-exchange 参数的值一致 return new FanoutExchange("exchange-unroute"); } /** * 申明死信交换器 * * @return */ @Bean public FanoutExchange dlxExchange() { return new FanoutExchange("exchange-dlx"); } /** * 申明绑定 * * @return */ @Bean public Binding binding() { return BindingBuilder.bind(queue()).to(exchange()).with("product").noargs(); } @Bean public Binding unRouteBinding() { return BindingBuilder.bind(unRouteQueue()).to(unRouteExchange()); } @Bean public Binding dlxBinding() { return BindingBuilder.bind(dlxQueue()).to(dlxExchange()); } }
注意:x-dead-letter-exchange和alternate-exchange参数的值和交换器中的值需要保持一致
5、生产者 - 编写消息发送者
@Component @Slf4j public class RabbitProducer implements ApplicationListener<ContextRefreshedEvent> { @Autowired private RabbitTemplate rabbitTemplate; @Override public void onApplicationEvent(ContextRefreshedEvent event) { String exchange = "exchange-rabbit-springboot-advance"; String routingKey = "product"; String unRoutingKey = "norProduct"; // 1.发送一条正常的消息 CorrelationData唯一(可以在ConfirmListener中确认消息) IntStream.rangeClosed(0, 10).forEach(num -> { String message = LocalDateTime.now().toString() + "发送第" + (num + 1) + "条消息."; rabbitTemplate.convertAndSend(exchange, routingKey, message, new CorrelationData("routing" + UUID.randomUUID().toString())); log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, routingKey, message); }); // 2.发送一条未被路由的消息,此消息将会进入备份交换器(alternate exchange) String message = LocalDateTime.now().toString() + "发送一条消息."; rabbitTemplate.convertAndSend(exchange, unRoutingKey, message, new CorrelationData("unRouting-" + UUID.randomUUID().toString())); log.info("发送一条消息,exchange:[{}],routingKey:[{}],message:[{}]", exchange, unRoutingKey, message); } }
注意:1、此处发送了2中消息,一种消息可以被正确的路由到消息队列中,另一种由于routingKey是不存在的,因此不会路由到队列中,观察这条消息有没有路由到 alternate-exchange 绑定的队列中。
2、CorrelationData数据需要唯一,此值可用于生产者确认消息。
6、生产者 - 启动类
@SpringBootApplication public class ProducerApplication { public static void main(String[] args) { SpringApplication.run(ProducerApplication.class, args); } }
7、消费者 - 配置文件
server: port: 9087 spring: rabbitmq: host: 140.143.237.224 port: 5672 username: root password: root virtual-host: / connection-timeout: 10000 listener: simple: acknowledge-mode: manual # 手动应答 auto-startup: true default-requeue-rejected: false # 不重回队列 concurrency: 5 max-concurrency: 20 prefetch: 1 # 每次只处理一个信息 retry: enabled: true
8、消费者 - 消息接收
@Component @Slf4j public class RabbitConsumer { /** * 监听 queue-rabbit-springboot-advance 队列 * * @param receiveMessage 接收到的消息 * @param message * @param channel */ @RabbitListener(queues = "queue-rabbit-springboot-advance") public void receiveMessage(String receiveMessage, Message message, Channel channel) { try { // 手动签收 log.info("接收到消息:[{}]", receiveMessage); if (new Random().nextInt(10) < 5) { log.warn("拒绝一条信息:[{}],此消息将会由死信交换器进行路由.", receiveMessage); channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, false); } else { channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } } catch (Exception e) { log.info("接收到消息之后的处理发生异常.", e); try { channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } catch (IOException e1) { log.error("签收异常.", e1); } } } }
注意:消费者接收中会随机拒绝几条消息,观察这个消息有没有进入 x-dead-letter-exchange 交换器绑定的队列中。
9、执行结果
完整代码:
代码如下:https://gitee.com/huan1993/rabbitmq/tree/master/rabbitmq-springboot-advanced
相关推荐
在SpringBoot中,配置一个Fanout交换器,不设置路由键,然后创建一个或多个队列绑定到这个交换器,即可实现广播效果。 3. **Direct模式**:直接模式基于路由键(Routing Key)匹配来决定消息发送到哪个队列。生产者...
- **创建DLX和关联队列**:首先,创建一个死信交换器,并设置业务队列的属性,使其在消息变为死信时将消息发布到DLX。 - **消息生产**:生产者发送消息到业务队列,可设置消息TTL。 - **消息处理**:消费者尝试...
1. **确认模式**:RabbitMQ支持发布者确认(Publisher Confirm),即消息发布后,RabbitMQ会向生产者返回一个确认,表示消息已成功路由到至少一个队列。Spring Boot可以开启这个特性,确保消息不会丢失。 2. **事务...
sprngboot整合 rabbitMQ MQ的作用 流量消峰 ,应用解耦,异步处理 ...1使用备份交换器路由到备胎队列消费。这样可以保证未被路由的消息不会丢失。 2通过消息的回调方法,添加ReturnListener的编程逻辑.
- 创建一个普通的队列,设置TTL属性,并将该队列绑定到前面创建的死信交换器。 - 当普通队列中的消息过期后,它会根据设置的死信路由键转发到死信队列,之后可以由消费者进行消费。 5. 代码实现 在Spring Boot中...
4. **Routing Key(路由键)**:在绑定中使用,帮助交换器确定消息应该被路由到哪个队列。路由键可以根据消息类型或其他特性进行设置。 在测试环境中,我们可以使用如下的步骤来验证生产者和消费者的交互: 1. **...
- 创建一个生产者类,使用`RabbitTemplate`发送消息到指定的交换器。 - 消息可以是简单的字符串,也可以是序列化的对象,取决于你的应用场景。 - 使用`convertAndSend`方法发送消息,其中可以指定交换器名称和...
1. 生产者发送消息到交换器,指定路由键(Routing Key)。 2. 交换器根据预设的规则(Binding)将消息路由到一个或多个队列。 3. 消费者从队列中接收消息,可以采用手动确认(acknowledgement)机制确保消息被正确...
本存储库`springboot-rabbitmq-retry-queues`主要关注的是如何在Spring Boot应用中实现RabbitMQ的非阻塞重试队列解决方案,这是一个针对处理消息处理失败时的策略,可以避免因单个任务失败而导致整个系统的阻塞。...
Spring Boot 提供了 @RabbitListener 注解来监听 RabbitMQ 队列中的消息,当处理消息时发生异常,可以通过配置消息确认(Acknowledgement)和死信队列(Dead Letter Queue)来实现重试。 1. **配置消息确认**: 设置...
- **RabbitMQ工作原理**:生产者发送消息到交换器,交换器根据预定义的路由规则将消息路由到一个或多个队列,消费者从队列中获取并消费消息。 2. **RabbitMQ安装与配置** - **安装流程**:涵盖在不同操作系统(如...
7. **实现实时接收**:在消费者端,可以设置消息确认模式为`Automatic Acknowledgement`,使得RabbitMQ在接收到消息后自动确认,实现实时消费。 以下是一个简单的生产者示例: ```java @Service public class ...
消息被发布到交换器,交换器根据预设的路由规则将消息分发到队列,队列再将消息路由给消费者。 在RabbitMQ实战中,你会学习到如何安装和配置RabbitMQ服务器,包括设置环境、启动和停止服务、管理用户、虚拟主机和...
4. **Routing Key(路由键)**:在发布/订阅模式中,生产者并不直接将消息发送到队列,而是通过交换器,路由键用来决定消息发往哪个具体的队列。 5. **Virtual Host(虚拟主机)**:虚拟主机是RabbitMQ中的逻辑隔离...
- **Exchange**:交换器是RabbitMQ的核心组件,它决定了消息如何路由到队列。常见的交换器类型有Direct、Fanout、Topic和Header。 - **Queue**:队列是存储消息的地方,每个消息至少被路由到一个队列。 - **...
1. **交换器(Exchanges)**:交换器负责接收生产者发送的消息,并根据预定义的路由规则将其分发到合适的队列。常见的交换器类型有Direct、Fanout、Topic和Header。 2. **队列(Queues)**:队列是消息的实际存储...
11. **Dead Letter Exchanges(死信交换器)**:当消息无法被正确处理时,可以配置死信交换器来处理这些消息,确保系统的稳定性。 通过这个项目,你可以深入理解RabbitMQ的基本工作原理,包括消息的发布与订阅、...
2. **交换器**(Exchange):接收来自生产者的消息后,根据配置规则将消息发送到指定的一个或多个队列。 3. **队列**(Queue):用来存储消息,直到被消费者消费。 4. **消费者**:接收并处理消息的一方。 #### 三...
生产者将消息发送到交换器,交换器根据预定义的路由规则将消息分发到一个或多个队列中。消费者从队列中获取并处理消息,这一过程可以是同步的,也可以是异步的。 **1. 交换器(Exchanges)** 交换器决定了消息如何...