前言
RabbitMQ 是基于 AMQP 协议的开源消息中间件,广泛应用于微服务架构中的异步解耦、削峰填谷和可靠投递场景。本文将从核心概念讲起,最终完成一个 Spring Boot 整合 RabbitMQ 的实战项目。
一、核心概念
1.1 AMQP 模型
Producer → Exchange → (Binding) → Queue → Consumer| 概念 | 说明 |
|---|---|
| Producer | 消息生产者 |
| Exchange | 交换机,接收消息并路由到队列 |
| Queue | 消息队列,存储消息 |
| Binding | 交换机与队列之间的绑定关系 |
| Routing Key | 路由键,Exchange 根据它决定消息去向 |
| Consumer | 消息消费者 |
1.2 四种交换机类型
| 类型 | 说明 | 适用场景 |
|---|---|---|
| Direct | 精确匹配 Routing Key | 点对点定向投递 |
| Fanout | 广播到所有绑定队列 | 广播通知 |
| Topic | 模式匹配 Routing Key | 按主题订阅 |
| Headers | 根据消息头匹配 | 复杂路由条件 |
二、Docker 快速部署
docker run -d --name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
-e RABBITMQ_DEFAULT_USER=admin \
-e RABBITMQ_DEFAULT_PASS=admin123 \
rabbitmq:3.13-management5672:AMQP 通信端口15672:管理界面端口,浏览器访问http://localhost:15672
三、Spring Boot 整合
3.1 添加依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>3.2 配置
spring:
rabbitmq:
host: localhost
port: 5672
username: admin
password: admin123
virtual-host: /3.3 队列与交换机配置
@Configuration
public class RabbitMQConfig {
public static final String EXCHANGE = "order.exchange";
public static final String QUEUE = "order.queue";
public static final String ROUTING_KEY = "order.create";
@Bean
public TopicExchange orderExchange() {
return new TopicExchange(EXCHANGE);
}
@Bean
public Queue orderQueue() {
return QueueBuilder.durable(QUEUE)
.withArgument("x-message-ttl", 60000)
.withArgument("x-dead-letter-exchange", "order.dlx")
.build();
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderExchange())
.with(ROUTING_KEY);
}
}3.4 生产者
@Service
public class OrderProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendOrder(OrderDTO order) {
rabbitTemplate.convertAndSend(
RabbitMQConfig.EXCHANGE,
RabbitMQConfig.ROUTING_KEY,
order,
message -> {
message.getMessageProperties()
.setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
);
}
}3.5 消费者
@Component
public class OrderConsumer {
@RabbitListener(queues = RabbitMQConfig.QUEUE)
public void handleOrder(OrderDTO order, Channel channel, Message message) throws IOException {
try {
processOrder(order);
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
channel.basicNack(
message.getMessageProperties().getDeliveryTag(),
false,
true
);
}
}
private void processOrder(OrderDTO order) {
System.out.println("处理订单: " + order.getId());
}
}3.6 配置手动 ACK
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual
prefetch: 10
retry:
enabled: true
max-attempts: 3四、可靠性投递方案
4.1 生产者确认(Publisher Confirm)
@Configuration
public class RabbitConfig {
@Bean
public RabbitTemplate rabbitTemplate(ConnectionFactory cf) {
RabbitTemplate template = new RabbitTemplate(cf);
template.setConfirmCallback((correlationData, ack, cause) -> {
if (!ack) {
log.error("消息未到达Exchange: {}", cause);
}
});
template.setReturnsCallback(returned -> {
log.error("消息未路由到Queue: {}", returned.getMessage());
});
template.setMandatory(true);
return template;
}
}4.2 死信队列(DLX)
消息变成死信的条件:
- 消息被消费者拒绝(basicNack + requeue=false)
- 消息 TTL 过期
- 队列达到最大长度
@Bean
public Queue deadLetterQueue() {
return QueueBuilder.durable("order.dlq").build();
}
@Bean
public FanoutExchange deadLetterExchange() {
return new FanoutExchange("order.dlx");
}
@Bean
public Binding dlxBinding() {
return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange());
}五、常见问题
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 消息丢失 | 消费者宕机未 ACK | 手动 ACK + 持久化 |
| 消息重复 | 网络抖动导致重复投递 | 消费端幂等性(唯一ID去重) |
| 消息堆积 | 消费速度小于生产速度 | 增加消费者 / 限流 |
| 队列阻塞 | 单条消息处理过慢 | 合理设置 prefetch |
总结
RabbitMQ 在微服务架构中承担着异步解耦和流量削峰的关键角色。生产环境中务必做到:持久化 + 手动 ACK + 死信队列 + 消费幂等,才能保证消息的可靠投递。
评论 (0)