RabbitMQ 集成
约 688 字大约 2 分钟
布欧-Lewyon
2026-05-17
首页 › Spring Boot › 消息与集成 › RabbitMQ 集成
RabbitMQ 是基于 AMQP 0-9-1 协议的消息代理(Message Broker),核心模型是 Exchange → Binding → Queue——消息不直接发给队列,而是通过 Exchange 根据 RoutingKey 路由。适合业务消息、异步解耦、延迟任务等场景。
快速入门
依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>配置
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: ${RABBIT_PASSWORD}
virtual-host: /
publisher-confirm-type: correlated # 生产者确认(CMP)
publisher-returns: true
listener:
simple:
acknowledge-mode: manual
prefetch: 10声明 Exchange/Queue/Binding
@Configuration
public class RabbitConfig {
@Bean
public DirectExchange userExchange() {
return new DirectExchange("user.exchange");
}
@Bean
public Queue userCreateQueue() {
return QueueBuilder.durable("user.create.queue")
.deadLetterExchange("user.dlx")
.ttl(60000)
.build();
}
@Bean
public Binding binding(Queue userCreateQueue, DirectExchange userExchange) {
return BindingBuilder.bind(userCreateQueue)
.to(userExchange).with("user.create");
}
}发送与消费
@Component
public class UserPublisher {
private final RabbitTemplate rabbitTemplate;
public void send(User user) {
rabbitTemplate.convertAndSend("user.exchange", "user.create", user);
}
}
@Component
public class UserConsumer {
@RabbitListener(queues = "user.create.queue")
public void handle(@Payload User user, Channel ch,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
try {
log.info("收到: {}", user);
ch.basicAck(tag, false);
} catch (Exception e) {
ch.basicNack(tag, false, true); // 重新入队
}
}
}Exchange 四种类型
| 类型 | 路由规则 | 声明方式 | 场景 |
|---|---|---|---|
| Direct | routingKey 精确匹配 | new DirectExchange("name") | 点对点命令 |
| Topic | * 匹配一级,# 匹配多级 | new TopicExchange("name") | 事件分发 |
| Fanout | 忽略 Key,广播到所有绑定 | new FanoutExchange("name") | 全局通知 |
| Headers | 按消息头 x-match 匹配 | new HeadersExchange("name") | 复杂路由条件 |
死信与重试
// 死信队列声明
@Bean
public Queue dlq() { return QueueBuilder.durable("user.dlx.queue").build(); }
@Bean
public DirectExchange dlx() { return new DirectExchange("user.dlx"); }
@Bean
public Binding dlqBinding() {
return BindingBuilder.bind(dlq()).to(dlx()).with("user.dead");
}生产确认
CorrelationData cd = new CorrelationData(UUID.randomUUID().toString());
rabbitTemplate.setConfirmCallback((correlation, ack, cause) -> {
if (!ack) {
log.error("投递失败: correlationId={}, cause={}", correlation.getId(), cause);
// 补偿:保存本地消息表 + 定时重扫
}
});
rabbitTemplate.setReturnsCallback(returned -> {
log.error("路由不可达: exchange={}, routingKey={}, replyText={}",
returned.getExchange(), returned.getRoutingKey(), returned.getReplyText());
});
rabbitTemplate.convertAndSend("user.exchange", "user.create", user, cd);架构:SimpleMessageListenerContainer
@EnableRabbit
└── RabbitListenerAnnotationBeanPostProcessor
└── 为每个 @RabbitListener 创建 MessageListenerContainer
└── SimpleMessageListenerContainer
├── 创建 Channel 连接
├── 创建 concurrency 个 Consumer 线程
├── 每个线程: basicConsume → receive → invoke listener
├── ACK 策略:AUTO(默认)/ MANUAL / NONE
└── 异常 → RetryTemplate 重试或 NACK 重新入队生产意识:publisher-confirm-type 在 Boot 3.x 必须为 correlated(旧版 publisher-confirms: true 已废弃)。死信队列必须配置 deadLetterExchange 和 deadLetterRoutingKey,否则 NACK 后消息直接丢弃。
小结
- RabbitMQ 通过 Exchange 路由消息到队列,四种 Exchange 覆盖点对点到广播场景。
- 手动 ACK + 死信 DLX + 生产确认 = 可靠投递三件套。
prefetch控制消费者预取数量,影响吞吐与公平性。
下一节:Kafka 集成 · RabbitMQ 与 Kafka 对比
上一节:生产部署
