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 协议的抽象和实现。
Spring Boot 提供了强大的支持,可以轻松地与 RabbitMQ 集成,简化了消息队列的使用。通过 Spring AMQP 项目,我们可以方便地配置、发送和接收消息,而无需编写大量的底层代码。
Spring AMQP (Advanced Message Queuing Protocol) 是 Spring 对 AMQP 协议的抽象和实现。它提供了一套高级的 API,用于与 AMQP 兼容的消息代理(如 RabbitMQ)进行交互。Spring AMQP 的核心组件包括:
RabbitTemplate: 用于发送消息。
MessageListenerContainer: 用于接收消息,通常使用 SimpleMessageListenerContainer 或 DirectMessageListenerContainer。
AmqpAdmin: 用于声明队列、交换机和绑定。
@RabbitListener: 注解,简化消息监听器的配置。
首先,需要在 pom.xml 文件中添加 Spring AMQP 的依赖:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>
在 application.properties 或 application.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
使用 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); }
使用 @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; } }
可以使用 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,并在应用程序启动时声明它们。
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); } }
RabbitMQ 提供了消息确认机制,以确保消息可靠地传递。Spring AMQP 支持以下两种确认模式:
NONE: 不进行确认。
AUTO: RabbitMQ 自动确认消息。
MANUAL: 手动确认消息。
要使用手动确认,需要在 application.properties 或 application.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(); } } }
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); // 执行数据库操作 } }
死信队列用于存储无法处理的消息。当消息被拒绝、过期或达到最大重试次数时,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。
以下是 RabbitMQ 与 Spring Boot 集成的基本流程图:
流程说明:
Spring Boot 应用程序使用 RabbitTemplate 发送消息。
RabbitTemplate 将消息发送到 RabbitMQ 交换机。
交换机根据路由规则将消息发送到相应的队列。
MessageListenerContainer 监听队列中的消息。
当队列中有消息时,MessageListenerContainer 将消息传递给消息监听器。
消息监听器处理消息。
Spring Boot 提供了强大的支持,可以轻松地与 RabbitMQ 集成。通过 Spring AMQP 项目,我们可以方便地配置、发送和接收消息,而无需编写大量的底层代码。本文介绍了 Spring Boot 与 RabbitMQ 集成的基本步骤,包括添加依赖、配置连接信息、发送消息、接收消息、声明队列、交换机和绑定、使用消息转换器、配置消息确认机制、使用事务和配置死信队列。希望本文能够帮助你更好地理解和使用 Spring Boot 与 RabbitMQ 集成。