本文目录导读:

RabbitMQ本身不直接支持延迟消息,但通过 延迟消息插件(rabbitmq_delayed_message_exchange) 可以实现消息的定时投递,以下是一个完整的技术实现方案:
插件安装与启用
# 1. 下载对应版本的插件 # 访问 https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases # 2. 将插件复制到RabbitMQ插件目录 cp rabbitmq_delayed_message_exchange-*.ez $RABBITMQ_HOME/plugins/ # 3. 启用插件 rabbitmq-plugins enable rabbitmq_delayed_message_exchange # 4. 确认插件状态 rabbitmq-plugins list | grep delayed
核心概念
延迟消息架构
生产者 → 延迟Exchange(x-delayed-type) → 绑定 → 队列 → 消费者
关键特性
- 使用
x-delayed-message类型的交换机 - 消息通过
x-delay头指定延迟时间(毫秒) - 支持多种路由模式:direct, topic, fanout, headers
Java 实现示例
Maven依赖
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.16.0</version>
</dependency>
生产者代码
import com.rabbitmq.client.*;
import java.util.HashMap;
import java.util.Map;
public class DelayedMessageProducer {
private static final String EXCHANGE_NAME = "delayed.exchange";
private static final String QUEUE_NAME = "delayed.queue";
private static final String ROUTING_KEY = "delayed.routing";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 1. 声明延迟交换机
Map<String, Object> argsMap = new HashMap<>();
argsMap.put("x-delayed-type", "direct");
channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message", true, false, argsMap);
// 2. 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 3. 绑定队列到交换机
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
// 4. 发送延迟消息
String message = "This is a delayed message";
// 设置消息属性
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.headers(new HashMap<String, Object>() {{
put("x-delay", 10000); // 10秒延迟
}})
.build();
channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, props, message.getBytes());
System.out.println(" [x] Sent delayed message: '" + message + "'");
}
}
}
消费者代码
import com.rabbitmq.client.*;
import java.util.HashMap;
import java.util.Map;
public class DelayedMessageConsumer {
private static final String EXCHANGE_NAME = "delayed.exchange";
private static final String QUEUE_NAME = "delayed.queue";
private static final String ROUTING_KEY = "delayed.routing";
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try (Connection connection = factory.newConnection();
Channel channel = connection.createChannel()) {
// 声明交换机(与生产者一致)
Map<String, Object> argsMap = new HashMap<>();
argsMap.put("x-delayed-type", "direct");
channel.exchangeDeclare(EXCHANGE_NAME, "x-delayed-message", true, false, argsMap);
// 声明队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 绑定
channel.queueBind(QUEUE_NAME, EXCHANGE_NAME, ROUTING_KEY);
System.out.println(" [*] Waiting for delayed messages. To exit press CTRL+C");
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), "UTF-8");
long deliveryTag = delivery.getEnvelope().getDeliveryTag();
System.out.println(" [x] Received message: '" + message + "'");
System.out.println(" [x] Delivered at: " + System.currentTimeMillis());
// 手动确认
channel.basicAck(deliveryTag, false);
};
channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> { });
// 保持连接
Thread.sleep(60000);
}
}
}
Spring Boot 集成示例
配置类
import org.springframework.amqp.core.*;
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 CustomExchange delayedExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange("delayed.exchange", "x-delayed-message", true, false, args);
}
@Bean
public Queue delayedQueue() {
return new Queue("delayed.queue", true);
}
@Bean
public Binding delayedBinding() {
return BindingBuilder
.bind(delayedQueue())
.to(delayedExchange())
.with("delayed.routing")
.noargs();
}
}
生产者Service
import org.springframework.amqp.AmqpException;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessagePostProcessor;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@Service
public class DelayedMessageService {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendDelayedMessage(String message, long delayMillis) {
MessagePostProcessor messagePostProcessor = new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
message.getMessageProperties().setDelay(Long.valueOf(delayMillis).intValue());
return message;
}
};
rabbitTemplate.convertAndSend(
"delayed.exchange",
"delayed.routing",
message,
messagePostProcessor
);
System.out.println("Sent delayed message: " + message + " with delay: " + delayMillis + "ms");
}
}
消费者监听器
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
@Component
public class DelayedMessageListener {
@RabbitListener(queues = "delayed.queue")
public void receiveDelayedMessage(String message) {
System.out.println("Received delayed message: " + message);
System.out.println("Received time: " + System.currentTimeMillis());
// 处理业务逻辑
}
}
高级特性与最佳实践
动态延迟时间
// 根据业务动态设置延迟
public void sendDynamicDelay(String message, int delaySeconds) {
MessagePostProcessor processor = msg -> {
msg.getMessageProperties().setDelay(delaySeconds * 1000);
return msg;
};
rabbitTemplate.convertAndSend(EXCHANGE, ROUTING_KEY, message, processor);
}
批量发送延迟消息
public void batchSendDelayedMessages(List<String> messages, int delaySeconds) {
for (String msg : messages) {
sendDynamicDelay(msg, delaySeconds);
}
}
延迟队列监控
@RestController
public class DelayMonitorController {
@Autowired
private RabbitManagementService managementService;
@GetMapping("/delayed-queue-status")
public Map<String, Object> getDelayedQueueStatus() {
Map<String, Object> status = new HashMap<>();
status.put("queueName", "delayed.queue");
status.put("messageCount", managementService.getMessageCount("delayed.queue"));
status.put("consumerCount", managementService.getConsumerCount("delayed.queue"));
return status;
}
}
注意事项
性能考虑
- 延迟消息存储在交换机级别的内存中
- 大量延迟消息可能影响性能
- 建议设置合理的消息过期时间
可靠性
// 开启发布确认 channel.confirmSelect(); // 开启事务(性能较低) channel.txSelect();
限制说明
- 最大延迟时间:约 2^32 毫秒(约49天)
- 不支持延迟消息的优先级
- 重启后延迟消息可能丢失(需配置持久化)
替代方案
如果插件不满足需求,可以考虑:
- 死信队列(DLX): 通过 TTL + 死信交换机实现
- Redis + 定时任务: 更灵活的延迟方案
- 时间轮算法: 高性能的延迟任务实现
死信队列方案(无插件)
@Configuration
public class DLXDelayConfig {
// 延迟队列
@Bean
public Queue delayQueue() {
Map<String, Object> args = new HashMap<>();
// 消息10秒后过期
args.put("x-message-ttl", 10000);
// 死信交换机
args.put("x-dead-letter-exchange", "process.exchange");
// 死信路由键
args.put("x-dead-letter-routing-key", "process.routing");
return new Queue("delay.queue", true, false, false, args);
}
// 实际处理队列
@Bean
public Queue processQueue() {
return new Queue("process.queue", true);
}
// 延迟交换机
@Bean
public DirectExchange delayExchange() {
return new DirectExchange("delay.exchange");
}
// 处理交换机
@Bean
public DirectExchange processExchange() {
return new DirectExchange("process.exchange");
}
@Bean
public Binding delayBinding() {
return BindingBuilder.bind(delayQueue()).to(delayExchange()).with("delay.routing");
}
@Bean
public Binding processBinding() {
return BindingBuilder.bind(processQueue()).to(processExchange()).with("process.routing");
}
}
这个方案提供了完整的 RabbitMQ 延迟消息实现,包括插件方案和无插件方案,可以根据实际需求选择使用。