RabbitMQ延迟插件定时投递

wen java案例 1

本文目录导读:

RabbitMQ延迟插件定时投递

  1. 插件安装与启用
  2. 核心概念
  3. Java 实现示例
  4. Spring Boot 集成示例
  5. 高级特性与最佳实践
  6. 注意事项
  7. 死信队列方案(无插件)

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 延迟消息实现,包括插件方案和无插件方案,可以根据实际需求选择使用。

上一篇Kafka生产者分区键决定分区

下一篇当前分类已是最新一篇

抱歉,评论功能暂时关闭!