RabbitMQ/Spring-AMQP 高级特性:事务机制与消息限流实战
在分布式系统中,消息的可靠投递至关重要。Spring AMQP 提供了多种机制来保障这一点,其中事务(Transaction)和限流(Flow Control)是生产环境中最常用的两个手段。
1. 事务机制
AMQP 协议实现了事务机制,主要用于确保消息发布的原子性。简单来说,它允许你将多个操作(如发送多条消息、确认操作)绑定在一起,要么全部成功,要么全部失败。
实现示例
在 Spring Boot 中,我们可以利用 @Transactional 注解配合配置好的 RabbitTemplate 来实现事务发送。
@RestController
@RequestMapping("/producer")
public class ProducerController {
@Resource(name = "transRabbitTemplate")
private RabbitTemplate transRabbitTemplate;
@Transactional
@RequestMapping("/trans")
public String trans() {
// 第一条消息
transRabbitTemplate.convertAndSend("", Constants.TRANS_QUEUE, "trans test ---> 1");
// 模拟业务异常,导致事务回滚
int num = 5 / 0;
// 这条消息不会发送
transRabbitTemplate.convertAndSend("", Constants.TRANS_QUEUE, "trans test ---> 2");
return "发送成功";
}
}
注意: Spring Boot 的 RabbitMQ 自动配置默认开启 Publisher Confirms 模式,但 RabbitMQ 不允许同一个通道同时使用事务模式和 Confirm 模式。因此,必须显式禁用 Confirm 功能。
spring:
rabbitmq:
publisher-confirm-type: none
publisher-returns: false
配置 Bean
我们需要手动创建支持事务的 RabbitTemplate 和事务管理器。
@Configuration
public class RabbitTemplateConfig {
@Bean("transRabbitTemplate")
public RabbitTemplate transRabbitTemplate(ConnectionFactory connectionFactory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
// 开启通道事务
rabbitTemplate.setChannelTransacted(true);
return rabbitTemplate;
}
@Bean
public RabbitTransactionManager rabbitTransactionManager(ConnectionFactory connectionFactory) {
return new RabbitTransactionManager(connectionFactory);
}
}
队列配置保持标准即可:
@Configuration
public class RabbitMQConfig {
@Bean("transQueue")
public Queue transQueue() {
return QueueBuilder.durable(Constants.TRANS_QUEUE).build();
}
}
2. 消息限流
消息限流(Flow Control)是防止生产者发送速度超过消费者处理能力的关键机制。如果处理不及时,会导致消息积压甚至系统崩溃。
生产者代码
这里演示一个简单的批量发送场景:
@RestController
@RequestMapping("/producer")
public class ProducerController {
@Resource(name = "rabbitTemplate")
private RabbitTemplate rabbitTemplate;
@RequestMapping("/qos")
public String qos() {
for (int i = 0; i < 20; i++) {
rabbitTemplate.convertAndSend(Constants.QOS_EXCHANGE, "qos", "qos test:" + i);
}
return "发送成功";
}
}
消费者配置与手动 ACK
为了实施限流,消费者端需要关闭自动确认,改为手动确认,并设置预取数量(Prefetch Count)。这意味着消费者每次只拉取少量消息,处理完一条再拉下一条。
@Component
@Slf4j
public class QosListener {
@RabbitListener(queues = Constants.QOS_QUEUE)
public void handMessage(Message message, Channel channel) throws IOException {
long deliveryTag = message.getMessageProperties().getDeliveryTag();
try {
log.info("接收到消息:{}, deliveryTag:{}",
new String(message.getBody(), StandardCharsets.UTF_8), deliveryTag);
log.info("处理成功");
// 处理成功后手动确认
// channel.basicAck(deliveryTag, true);
} catch (Exception e) {
// 处理失败,拒绝消息并重新入队
channel.basicNack(deliveryTag, true, true);
}
}
}
关键配置
在 application.yml 中,我们需要将监听器模式设为 manual,并限制 prefetch 数量。
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 开启手动确认
prefetch: 5 # 每个消费者最多未确认 5 条消息
最后,确保交换器、队列和绑定关系配置正确:
@Configuration
public class RabbitMQConfig {
@Bean("qosQueue")
public Queue qosQueue() {
return QueueBuilder.durable(Constants.QOS_QUEUE).build();
}
@Bean("qosExchange")
public DirectExchange qosExchange() {
return ExchangeBuilder.directExchange(Constants.QOS_EXCHANGE).build();
}
@Bean("qosBinding")
public Binding qosBinding(@Qualifier("qosExchange") DirectExchange directExchange,
@Qualifier("qosQueue") Queue queue) {
return BindingBuilder.bind(queue).to(directExchange).with("qos");
}
}
通过合理配置 Prefetch 和手动 ACK,我们可以在保证消息不丢失的前提下,有效控制消费端的压力,避免系统过载。

