5.5 RabbitMQ 与 Spring Boot 集成


文档摘要

5.5 RabbitMQ 与 Spring Boot 集成 (Integration with Spring Boot) 5.5 RabbitMQ 与 Spring Boot 集成 (Integration with Spring Boot) Spring Boot 提供了强大的支持,可以轻松地与 RabbitMQ 集成,简化了消息队列的使用。通过 Spring AMQP 项目,我们可以方便地配置、发送和接收消息,而无需编写大量的底层代码。 5.5.1 Spring AMQP 简介 Spring AMQP (Advanced Message Queuing Protocol) 是 Spring 对 AMQP 协议的抽象和实现。

5.5 RabbitMQ 与 Spring Boot 集成 (Integration with Spring Boot)

5.5 RabbitMQ 与 Spring Boot 集成 (Integration with Spring Boot)

Spring Boot 提供了强大的支持,可以轻松地与 RabbitMQ 集成,简化了消息队列的使用。通过 Spring AMQP 项目,我们可以方便地配置、发送和接收消息,而无需编写大量的底层代码。

5.5.1 Spring AMQP 简介

Spring AMQP (Advanced Message Queuing Protocol) 是 Spring 对 AMQP 协议的抽象和实现。它提供了一套高级的 API,用于与 AMQP 兼容的消息代理(如 RabbitMQ)进行交互。Spring AMQP 的核心组件包括:

  • RabbitTemplate: 用于发送消息。

  • MessageListenerContainer: 用于接收消息,通常使用 SimpleMessageListenerContainerDirectMessageListenerContainer

  • AmqpAdmin: 用于声明队列、交换机和绑定。

  • @RabbitListener: 注解,简化消息监听器的配置。

5.5.2 添加依赖

首先,需要在 pom.xml 文件中添加 Spring AMQP 的依赖:

<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>

5.5.3 配置 RabbitMQ 连接

application.propertiesapplication.yml 文件中配置 RabbitMQ 连接信息:

spring.rabbitmq.host=localhost spring.rabbitmq.port=5672 spring.rabbitmq.username=guest spring.rabbitmq.password=guest

或者使用 YAML 格式:

spring: rabbitmq: host: localhost port: 5672 username: guest password: guest

5.5.4 发送消息

使用 RabbitTemplate 发送消息。首先,注入 RabbitTemplate

import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; @Component public class MessageSender { @Autowired private RabbitTemplate rabbitTemplate; public void sendMessage(String exchange, String routingKey, String message) { rabbitTemplate.convertAndSend(exchange, routingKey, message); } }

然后,在需要发送消息的地方调用 sendMessage 方法:

@Autowired private MessageSender messageSender; public void sendNotification(String message) { messageSender.sendMessage("my.exchange", "my.routing.key", message); }

5.5.5 接收消息

使用 @RabbitListener 注解或 MessageListenerContainer 接收消息。

使用 @RabbitListener 注解:

import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class MessageReceiver { @RabbitListener(queues = "my.queue") public void receiveMessage(String message) { System.out.println("Received message: " + message); } }

使用 MessageListenerContainer

import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitConfig { @Bean public Queue myQueue() { return new Queue("my.queue", false); } @Bean public SimpleMessageListenerContainer container(ConnectionFactory connectionFactory) { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setQueueNames("my.queue"); container.setMessageListener(message -> { String messageBody = new String(message.getBody()); System.out.println("Received message: " + messageBody); }); return container; } }

5.5.6 声明队列、交换机和绑定

可以使用 AmqpAdmin 或 Spring Bean 来声明队列、交换机和绑定。

使用 AmqpAdmin

import org.springframework.amqp.core.AmqpAdmin; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import javax.annotation.PostConstruct; @Configuration public class RabbitConfig { @Autowired private AmqpAdmin amqpAdmin; @Bean public Queue myQueue() { return new Queue("my.queue", false); } @Bean public DirectExchange myExchange() { return new DirectExchange("my.exchange"); } @Bean public Binding myBinding(Queue myQueue, DirectExchange myExchange) { return BindingBuilder.bind(myQueue).to(myExchange).with("my.routing.key"); } @PostConstruct public void createQueuesAndBindings() { amqpAdmin.declareQueue(myQueue()); amqpAdmin.declareExchange(myExchange()); amqpAdmin.declareBinding(myBinding(myQueue(), myExchange())); } }

使用 Spring Bean:

import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitConfig { @Bean public Queue myQueue() { return new Queue("my.queue", false); } @Bean public DirectExchange myExchange() { return new DirectExchange("my.exchange"); } @Bean public Binding myBinding(Queue myQueue, DirectExchange myExchange) { return BindingBuilder.bind(myQueue).to(myExchange).with("my.routing.key"); } }

Spring Boot 会自动检测这些 Bean,并在应用程序启动时声明它们。

5.5.7 消息转换器

Spring AMQP 提供了多种消息转换器,用于将消息内容转换为 Java 对象,以及将 Java 对象转换为消息内容。常用的消息转换器包括:

  • SimpleMessageConverter: 默认的转换器,使用 Java 序列化。

  • Jackson2JsonMessageConverter: 使用 Jackson 库进行 JSON 序列化和反序列化。

  • TextMessageConverter: 用于处理文本消息。

要使用 Jackson2JsonMessageConverter,需要添加 Jackson 依赖:

<dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> </dependency>

然后,配置 RabbitTemplate 使用 Jackson2JsonMessageConverter

import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitConfig { @Bean public RabbitTemplate rabbitTemplate(org.springframework.amqp.rabbit.connection.ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter()); return rabbitTemplate; } @Bean public Jackson2JsonMessageConverter jsonMessageConverter() { return new Jackson2JsonMessageConverter(); } }

现在,可以发送和接收 Java 对象:

@Component public class MessageSender { @Autowired private RabbitTemplate rabbitTemplate; public void sendMessage(String exchange, String routingKey, Object message) { rabbitTemplate.convertAndSend(exchange, routingKey, message); } } @Component public class MessageReceiver { @RabbitListener(queues = "my.queue") public void receiveMessage(MyObject message) { System.out.println("Received message: " + message); } }

5.5.8 消息确认机制

RabbitMQ 提供了消息确认机制,以确保消息可靠地传递。Spring AMQP 支持以下两种确认模式:

  • NONE: 不进行确认。

  • AUTO: RabbitMQ 自动确认消息。

  • MANUAL: 手动确认消息。

要使用手动确认,需要在 application.propertiesapplication.yml 文件中配置:

spring.rabbitmq.listener.simple.acknowledge-mode=manual

然后,在消息监听器中手动确认消息:

import com.rabbitmq.client.Channel; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; import java.io.IOException; @Component public class MessageReceiver { @RabbitListener(queues = "my.queue") public void receiveMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) { System.out.println("Received message: " + message); try { channel.basicAck(tag, false); } catch (IOException e) { // 处理确认失败的情况 e.printStackTrace(); } } }

5.5.9 事务

Spring AMQP 支持事务,可以确保消息发送和数据库操作的原子性。要使用事务,需要配置 RabbitTemplate 使用 ChannelTransacted

import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.transaction.annotation.EnableTransactionManagement; @Configuration @EnableTransactionManagement public class RabbitConfig { @Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); rabbitTemplate.setChannelTransacted(true); return rabbitTemplate; } }

然后,在需要事务的方法上添加 @Transactional 注解:

import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @Service public class MyService { @Autowired private MessageSender messageSender; @Transactional public void sendMessageAndSaveData(String message) { messageSender.sendMessage("my.exchange", "my.routing.key", message); // 执行数据库操作 } }

5.5.10 死信队列 (Dead Letter Queue, DLQ)

死信队列用于存储无法处理的消息。当消息被拒绝、过期或达到最大重试次数时,RabbitMQ 会将消息发送到死信队列。

首先,声明死信队列和交换机:

import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; @Configuration public class RabbitConfig { @Bean public Queue dlq() { return new Queue("dlq", false); } @Bean public DirectExchange dlx() { return new DirectExchange("dlx"); } @Bean public Binding dlqBinding(Queue dlq, DirectExchange dlx) { return BindingBuilder.bind(dlq).to(dlx).with("dlq.routing.key"); } @Bean public Queue myQueue() { Map<String, Object> args = new HashMap<>(); args.put("x-dead-letter-exchange", "dlx"); args.put("x-dead-letter-routing-key", "dlq.routing.key"); return new Queue("my.queue", false, false, false, args); } @Bean public DirectExchange myExchange() { return new DirectExchange("my.exchange"); } @Bean public Binding myBinding(Queue myQueue, DirectExchange myExchange) { return BindingBuilder.bind(myQueue).to(myExchange).with("my.routing.key"); } }

然后,配置消息监听器,当消息处理失败时,拒绝消息:

import com.rabbitmq.client.Channel; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; import java.io.IOException; @Component public class MessageReceiver { @RabbitListener(queues = "my.queue") public void receiveMessage(String message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) { System.out.println("Received message: " + message); try { // 模拟处理失败 throw new RuntimeException("Processing failed"); } catch (Exception e) { try { channel.basicNack(tag, false, false); // 拒绝消息,并将其发送到死信队列 } catch (IOException ex) { ex.printStackTrace(); } } } }

现在,当消息处理失败时,它将被发送到死信队列 dlq

5.5.11 集成流程图

以下是 RabbitMQ 与 Spring Boot 集成的基本流程图:

流程说明:

  1. Spring Boot 应用程序使用 RabbitTemplate 发送消息。

  2. RabbitTemplate 将消息发送到 RabbitMQ 交换机。

  3. 交换机根据路由规则将消息发送到相应的队列。

  4. MessageListenerContainer 监听队列中的消息。

  5. 当队列中有消息时,MessageListenerContainer 将消息传递给消息监听器。

  6. 消息监听器处理消息。

5.5.12 总结

Spring Boot 提供了强大的支持,可以轻松地与 RabbitMQ 集成。通过 Spring AMQP 项目,我们可以方便地配置、发送和接收消息,而无需编写大量的底层代码。本文介绍了 Spring Boot 与 RabbitMQ 集成的基本步骤,包括添加依赖、配置连接信息、发送消息、接收消息、声明队列、交换机和绑定、使用消息转换器、配置消息确认机制、使用事务和配置死信队列。希望本文能够帮助你更好地理解和使用 Spring Boot 与 RabbitMQ 集成。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U