RabbitMQ + Spring Boot 消息队列
Spring Boot 整合 RabbitMQ 完整示例,含声明、发送、消费和确认。
RabbitMQ + Spring Boot 消息队列实战
RabbitMQ 核心模型
RabbitMQ 基于 AMQP 协议,消息流路径:Producer → Exchange → Queue → Consumer。
与 Kafka 不同,RabbitMQ 的 Exchange 概念让消息路由非常灵活。Producer 不直接发消息到队列,而是发到 Exchange,Exchange 根据绑定规则决定消息投递到哪个队列。
四种交换机
1. Direct Exchange(精确匹配)
路由键精确匹配。比如订单服务发送 order.created 消息,只有绑定了 order.created 路由键的队列能收到:
@Bean
public Binding binding(Queue queue, DirectExchange exchange) {
return BindingBuilder.bind(queue).to(exchange).with("order.created");
}
2. Topic Exchange(模式匹配)
路由键支持通配符:# 匹配零或多个单词,* 匹配一个单词。适合多级分类:
// 消费者监听所有 order 相关消息
@RabbitListener(bindings = @QueueBinding(
exchange = @Exchange(name = "order.topic", type = "topic"),
value = @Queue,
key = "order.#"
))
public void handleOrder(OrderMessage msg) { ... }
3. Fanout Exchange(广播)
消息发送到所有绑定的队列,忽略路由键。适合全局通知、缓存刷新:
@Bean
public FanoutExchange cacheRefresh() {
return new FanoutExchange("cache.refresh");
}
4. Headers Exchange(Header 匹配)
不依赖路由键,根据消息 Header 属性匹配。灵活性最高但性能较差,实际使用较少。
Spring Boot 整合
Spring Boot 通过 spring-boot-starter-amqp 自动配置:
spring:
rabbitmq:
host: 192.168.8.133
port: 5672
username: admin
password: admin123
virtual-host: /
listener:
simple:
retry:
enabled: true
max-attempts: 3
initial-interval: 2000
我的 VM 上使用 Docker 部署 RabbitMQ 3:
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 -e RABBITMQ_DEFAULT_USER=admin -e RABBITMQ_DEFAULT_PASS=admin123 rabbitmq:3-management
管理控制台通过 http://192.168.8.133:15672 访问,可以查看队列的堆积情况、消费者状态、消息速率等。
实际应用场景
在 CheckByAI 项目中,RabbitMQ 作为任务分发的核心组件:
// 生产者:提交审核任务
rabbitTemplate.convertAndSend("check.exchange", "task.submit", taskMessage);
// 消费者:处理审核任务 + 手动 ACK
@RabbitListener(queues = "check.task.queue")
public void handleTask(TaskMessage task, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
processTask(task);
channel.basicAck(tag, false); // 处理成功才 ACK
} catch (Exception e) {
channel.basicNack(tag, false, true); // 重试
}
}
核心价值:
- 异步解耦:任务提交方不需要等待处理完成
- 削峰填谷:瞬时流量高峰消息积在队列中慢慢消费
- 死信队列:重试多次失败的消息进入 DLQ 人工处理
生产踩坑
消息堆积告警:某个消费者宕机后消息在队列中持续堆积,因为没设置队列长度上限和告警阈值。解决方案:配置队列的 x-max-length 和监控告警。
手动 ACK vs 自动 ACK:自动 ACK 如果消费者处理中崩溃,消息就丢了。生产环境必须使用手动 ACK + 重试机制。
幂等消费:RabbitMQ 的 at-least-once 语义意味着同一条消息可能被消费多次。消费端需要利用业务唯一 ID 做幂等判断,防止重复处理。