因为spring boot 1.4版本比较老了,所以不能直接支持延时队列。所以需要用手工创建Exchange

    @RabbitListener(
            queues = "test-log-queue"
    )
    public void processMessage(String message) {
        logger.info("我接收了[{}]。", message);
    }

    @Bean
    public ApplicationRunner runner(@Qualifier("amqpTemplate") RabbitTemplate template) {
        return args -> {
            MessageProperties messageProperties = new MessageProperties();
            template.convertAndSend("test-log-exchange", "main", "34020000001320000001", message->{
                message.getMessageProperties().setHeader("x-delay", 60 * 1000);
                return message;
            });
            logger.info("我发送成功了!");
        };
    }
    @Bean
    public org.springframework.amqp.core.Queue manualDelayQueue() {
        // 持久化、非排他、非自动删除
        return new org.springframework.amqp.core.Queue("test-log-queue", true, false, false);
    }

    // 3. 手动声明交换机 + 绑定队列(核心:PostConstruct初始化)
    @PostConstruct
    public void declareManualExchangeAndBinding() {
        // 3.1 构建自定义延迟交换机(x-delayed-message类型)
        Map<String, Object> exchangeArgs = new HashMap<>();
        exchangeArgs.put("x-delayed-type", "direct"); // 指定延迟交换机的路由类型
        Map<String, Object> args = new HashMap<>();
        args.put("x-delayed-type", "direct"); // 指定内部交换机类型为 direct
        org.springframework.amqp.core.Exchange manualDelayExchange = new CustomExchange("test-log-exchange", "x-delayed-message", true, false, args);
        // 3.2 手动创建交换机(绕过Spring AMQP的类型校验)
        amqpAdmin.declareExchange(manualDelayExchange);

        // 3.3 手动绑定队列到交换机(指定路由键)
        Binding binding = BindingBuilder.bind(manualDelayQueue())
                .to(manualDelayExchange)
                .with("main")
                .noargs();
        amqpAdmin.declareBinding(binding);

       logger.info("手动声明交换机&绑定完成:" + "test-log-exchange");
    }

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐