基于 Spring Boot 的 RabbitMQ 核心机制与生产实践
Queue(队列)是 RabbitMQ 中保存消息的容器,消息必须投递到队列里,消费者再从队列中拉取。RabbitMQ 与 Kafka 的最大区别在于:RabbitMQ 以队列为核心存储单元,而 Kafka 以 Topic 为逻辑层面、以分区文件为物理存储。
当多个消费者订阅同一个队列时,RabbitMQ 会采用轮询(Round-Robin)方式把消息分摊给这些消费者,因此一个消息不会被所有消费者同时消费。需要注意的是,RabbitMQ 原生不支持"队列级广播",如果需要广播,应使用 Fanout 交换机绑定多个队列来实现,而不是在队列层面做二次开发。
交换机(Exchange)常见类型有四种:
- Direct:精确匹配路由键。
- Fanout:广播到所有绑定队列,忽略路由键。
- Topic:支持通配符
*和#匹配路由键。 - Headers:根据消息头属性匹配。
2. Spring Boot 基础配置
下面给出一个完整的 application.yml 示例,涵盖集群地址、发布确认、Return 回调、Mandatory 设置以及消费者手动确认等关键项。
server:
port: 11000
spring:
application:
name: rabbitmq-demo
rabbitmq:
addresses: node1:5672,node2:5672
username: admin
password: admin
virtual-host: /demo
publisher-returns: true
publisher-confirm-type: correlated
connection-timeout: 15000
template:
mandatory: true
default-exchange: demo.direct
listener:
simple:
acknowledge-mode: manual
concurrency: 5
max-concurrency: 10
prefetch: 1
retry:
enabled: true
关键配置说明:
publisher-returns/publisher-confirm-type:开启 Return 与 Confirm 两种回调。template.mandatory:交换机无法路由到队列时,消息是否返回生产者,必须设置为true才能触发 ReturnCallback。acknowledge-mode: manual:消费端手动确认,业务处理完成后再 ACK。prefetch:限制每个消费者同时只能处理指定条数未确认消息,常用于削峰。
3. Java 配置类
通过 @Configuration 声明交换机、队列、绑定关系,并创建两个 RabbitTemplate 示例:一个使用自定义连接工厂,一个直接注入 Spring 管理的连接工厂。
@Configuration
public class AmqpConfig {
@Value("${spring.rabbitmq.addresses}")
private String addresses;
@Value("${spring.rabbitmq.username}")
private String user;
@Value("${spring.rabbitmq.password}")
private String pwd;
@Value("${spring.rabbitmq.virtual-host}")
private String vhost;
@Value("${spring.rabbitmq.default-exchange}")
private String defaultExchange;
@Autowired
private ConfirmCallbackService confirmCallback;
@Autowired
private ReturnCallbackService returnCallback;
@Bean
public ConnectionFactory amqpConnectionFactory() {
CachingConnectionFactory factory = new CachingConnectionFactory();
factory.setAddresses(addresses);
factory.setUsername(user);
factory.setPassword(pwd);
factory.setVirtualHost(vhost);
return factory;
}
@Bean("customRabbitTemplate")
public RabbitTemplate customRabbitTemplate(ConnectionFactory amqpConnectionFactory) {
RabbitTemplate template = new RabbitTemplate(amqpConnectionFactory);
template.setMessageConverter(jsonMessageConverter());
template.setExchange(defaultExchange);
template.setMandatory(true);
template.setConfirmCallback(confirmCallback);
template.setReturnsCallback(returnCallback);
template.setReplyTimeout(20000);
template.setReceiveTimeout(20000);
return template;
}
@Bean
public Jackson2JsonMessageConverter jsonMessageConverter() {
return new Jackson2JsonMessageConverter();
}
@Bean
public DirectExchange orderDirectExchange() {
return new DirectExchange("order.direct", true, false);
}
@Bean
public FanoutExchange logFanoutExchange() {
return new FanoutExchange("log.fanout", true, false);
}
@Bean
public TopicExchange noticeTopicExchange() {
return new TopicExchange("notice.topic", true, false);
}
@Bean
public Queue orderQueue() {
return new Queue("order.queue", true);
}
@Bean
public Queue smsQueue() {
return new Queue("sms.queue", true);
}
@Bean
public Queue emailQueue() {
return new Queue("email.queue", true);
}
@Bean
public Binding orderBinding() {
return BindingBuilder.bind(orderQueue())
.to(orderDirectExchange())
.with("order.created");
}
@Bean
public Binding smsBinding() {
return BindingBuilder.bind(smsQueue()).to(logFanoutExchange());
}
@Bean
public Binding emailBinding() {
return BindingBuilder.bind(emailQueue())
.to(noticeTopicExchange())
.with("notice.#");
}
}
声明交换机时主要关注三个属性:durable(是否持久化)、autoDelete(没有队列绑定时是否自动删除)、以及自定义扩展参数 arguments。
4. 消息转换器
RabbitTemplate 通过 MessageConverter 完成 Java 对象与 Message 之间的转换。可以自定义转换器,也可以直接使用 Jackson2JsonMessageConverter。
public class PlainTextConverter implements MessageConverter {
@Override
public Message toMessage(Object object, MessageProperties properties) throws MessageConversionException {
return new Message(object.toString().getBytes(StandardCharsets.UTF_8), properties);
}
@Override
public Object fromMessage(Message message) throws MessageConversionException {
return new String(message.getBody(), StandardCharsets.UTF_8);
}
}
使用 Jackson2JsonMessageConverter 时,发送端与接收端的类型必须一致,包括类名、字段以及包路径,否则反序列化会失败。如果业务中需要多态或松耦合,建议自定义转换器或在消息头里存放类型信息。
5. 生产者与消费者示例
下面展示一个统一封装的发送服务,通过 MessagePostProcessor 在发送前设置消息 ID、Correlation ID 以及持久化投递模式。
@Service
@Slf4j
public class MsgPublisher {
@Resource(name = "customRabbitTemplate")
private RabbitTemplate rabbitTemplate;
public void sendObject(String routingKey, Object payload) {
rabbitTemplate.convertAndSend(routingKey, payload, this::decorateMessage);
}
public void sendObject(String exchange, String routingKey, Object payload) {
rabbitTemplate.convertAndSend(exchange, routingKey, payload,
this::decorateMessage,
new CorrelationData(UUID.randomUUID().toString()));
}
private Message decorateMessage(Message message) {
message.getMessageProperties().setMessageId(UUID.randomUUID().toString());
message.getMessageProperties().setCorrelationId(UUID.randomUUID().toString());
message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return message;
}
}
消费端可以使用 @RabbitListener 标注在类上,再搭配多个 @RabbitHandler 方法,Spring 会根据 MessageConverter 转换后的实际类型分派到对应方法。
@Component
@Slf4j
@RabbitListener(queues = "notice.queue")
public class NoticeConsumer {
@RabbitHandler
public void onString(@Payload String text) {
log.info("收到文本消息:{}", text);
}
@RabbitHandler
public void onOrder(@Payload Order order) {
log.info("收到订单消息:{}", order);
}
@RabbitHandler
public void onMap(@Payload Map map) {
map.forEach((k, v) -> log.info("{}={}", k, v));
}
}
需要注意:如果多个 @RabbitHandler 方法都接受 Map 类型,Spring 无法确定调用哪个方法,会抛出异常。因此不建议把 Map 作为通用接收类型。
6. 消息确认与手动 ACK
生产端需要实现 RabbitTemplate.ConfirmCallback 和 RabbitTemplate.ReturnsCallback,分别处理"是否送达交换机"和"是否能路由到队列"。
@Component
@Slf4j
public class OrderConfirmCallback implements RabbitTemplate.ConfirmCallback {
@Override
public void confirm(CorrelationData data, boolean ack, String cause) {
if (ack) {
log.info("消息已到达交换机,id={}", data != null ? data.getId() : "");
} else {
log.error("消息未到达交换机,id={},原因:{}", data != null ? data.getId() : "", cause);
}
}
}
@Component
@Slf4j
public class OrderReturnCallback implements RabbitTemplate.ReturnsCallback {
@Override
public void returnedMessage(ReturnedMessage returned) {
log.error("消息被退回,exchange={},routingKey={},replyCode={},replyText={}",
returned.getExchange(),
returned.getRoutingKey(),
returned.getReplyCode(),
returned.getReplyText());
}
}
消费端开启手动确认后,需要在业务成功后主动 ACK;失败时可以选择重新入队或拒绝。
@RabbitHandler
public void handleOrder(@Payload Order order, Channel channel, Message message) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
// 业务处理
channel.basicAck(tag, false);
} catch (Exception e) {
if (message.getMessageProperties().getRedelivered()) {
channel.basicReject(tag, false);
} else {
channel.basicNack(tag, false, true);
}
}
}
三种回执方法的区别如下:
basicAck(deliveryTag, multiple):确认成功,multiple=true时批量确认小于当前deliveryTag的所有消息。basicNack(deliveryTag, multiple, requeue):确认失败,requeue=true时重新入队。basicReject(deliveryTag, requeue):单条拒绝,不能批量。
四种投递场景对应的回调触发情况:
- 找不到交换机:
ConfirmCallback返回 false。 - 找到交换机但找不到队列:
ConfirmCallback返回 true,同时触发ReturnCallback。 - 交换机和队列都找不到:同第一种,
ConfirmCallback返回 false。 - 投递成功:
ConfirmCallback返回 true。