如何使用Spring Boot防止RabbitMQ队列中的消息被自动删除?
Spring Boot 实现RabbitMQ消息自定义留存时长的方案
默认情况下RabbitMQ使用自动确认(autoAck=true)模式,消费者拿到消息后Broker就会直接标记消息为可删除,因此无法自定义控制消息的留存时间,Spring Boot生态完全支持该需求的实现,核心通过调整消费确认模式+消息TTL(过期时间)配置完成。
1. 基础依赖引入
项目中先引入RabbitMQ对应的Starter依赖,已引入可跳过该步:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>
2. 配置手动确认模式
先关闭默认的自动消费确认,避免消息被自动删除,在application.yml中添加如下配置:
spring: rabbitmq: host: 你的RabbitMQ服务地址 port: 5672 username: 账号 password: 密码 listener: simple: acknowledge-mode: manual # 开启手动确认,消费后消息不会自动删除 default-requeue-rejected: true # 可选配置,消费失败的消息默认放回队列
3. 自定义消息留存时长的两种实现方式
队列全局统一留存时长:给整个队列的所有消息设置统一的过期时间,到期后消息自动被删除
队列配置代码示例:import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; @Configuration public class RabbitMqConfig { @Bean public Queue customRetentionQueue() { Map<String, Object> args = new HashMap<>(); // 单位为毫秒,示例设置消息统一留存1小时 args.put("x-message-ttl", 3600000); return new Queue("custom-retention-queue", true, false, false, args); } }单条消息独立设置留存时长:如果不同消息需要不同的留存时间,可在发送消息时单独指定过期时间
消息发送代码示例:import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageProperties; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; import javax.annotation.Resource; @Component public class MessageSender { @Resource private RabbitTemplate rabbitTemplate; public void sendMsgWithCustomTtl(String content, long ttlMillis) { MessageProperties properties = new MessageProperties(); // 单独设置当前消息的留存时长,单位为毫秒 properties.setExpiration(String.valueOf(ttlMillis)); Message message = new Message(content.getBytes(), properties); rabbitTemplate.convertAndSend("custom-retention-queue", message); } }
4. 手动控制消息删除时机
如果不需要等TTL到期自动删除,可在业务处理完成后手动确认删除消息,消费端代码示例:
import com.rabbitmq.client.Channel; import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; import java.io.IOException; @Component public class CustomMessageConsumer { @RabbitListener(queues = "custom-retention-queue") public void consume(Message message, Channel channel) throws IOException { long deliveryTag = message.getMessageProperties().getDeliveryTag(); try { // 你的业务处理逻辑 String msgContent = new String(message.getBody()); // 业务处理完成后手动确认,此时消息才会从队列中删除 channel.basicAck(deliveryTag, false); // 若业务处理失败需要将消息放回队列重试,可调用: // channel.basicNack(deliveryTag, false, true); // 若业务处理失败不需要保留消息,可调用: // channel.basicNack(deliveryTag, false, false); } catch (Exception e) { // 异常场景根据业务需求决定消息是否放回队列 channel.basicNack(deliveryTag, false, true); } } }
注意事项
- 若同时配置了队列全局TTL和单条消息TTL,RabbitMQ会取两者中更小的值作为实际的消息过期时间
- 队列的
x-message-ttl参数只能在队列首次创建时设置,已存在的队列需要先删除再重新创建才能修改该参数 - 手动确认模式下必须在逻辑处理完成后调用ack/nack方法,否则消息会一直留在队列中,直到消费者断开连接才会重新被放回队列
内容的提问来源于stack exchange,提问作者Maksudur Rahman Maruf
相关产品推荐
相关产品推荐

