Issue with RabbitMQ Delay message in spring boot

Viewed 22

I am facing an issue in Rabbit MQ regarding x-delay while connecting to spring boot. I need to schedule the messages for a variable delay according to the message type. It can be one of the units MINUTE, DAY, WEEK, MONTH, and so on…

Below is my configuration class :

private final RabbitProperties rabbitProperties;
    private final Environment environment;

    @Bean
    public Queue rabbitMQueue() {
        return new Queue(environment.getProperty(RABBITMQ_QUEUE_NAME), false);
    }

    @Bean
    public Exchange rabbitExchange() {

        String exchangeName = "test_exchange";

        Map<String, Object> exchangeArgs = new HashMap<>();
        exchangeArgs.put("x-delayed-type", exchangeType.toLowerCase());
        exchangeArgs.put("x-delayed-message",true);
        exchangeArgs.put("x-message-ttl",9922);

        log.info("Loading {} exchange with name {}.", exchangeType, exchangeName);

        switch (exchangeType){
            default: return new CustomExchange(exchangeName, exchangeType, true, false, exchangeArgs);
            case "DIRECT" : return directExchange(exchangeName, exchangeArgs);
        }
    }

    private Exchange directExchange(String exchangeName, Map<String, Object> exchangeArgs) {
//        log.info("Generating directExchange");
//        DirectExchange directExchange = new DirectExchange(exchangeName,true, false, exchangeArgs);
//        directExchange.setDelayed(true);
//        return directExchange;

        return ExchangeBuilder.directExchange(exchangeName).withArguments(exchangeArgs)
                .delayed()
                .build();
    }

    @Bean
    public Binding rabbitBinding(final Queue rabbitMQueue, final Exchange rabbitExchange){

        log.info("Exchange to bind : {}", rabbitExchange.getName());
        return BindingBuilder
                .bind(rabbitMQueue)
                .to(rabbitExchange)
                .with(environment.getProperty(RABBITMQ_ROUTING_KEY)).noargs();
    }

    @Bean
    public AmqpTemplate amqpTemplate(final ConnectionFactory rabbitMQConnectionFactory,
                                     final MessageConverter rabbitMessageConvertor,
                                     final Exchange rabbitExchange){

        RabbitTemplate rabbitTemplate = new RabbitTemplate(rabbitMQConnectionFactory);
        rabbitTemplate.setMessageConverter(rabbitMessageConvertor);
        rabbitTemplate.setExchange(rabbitExchange.getName());
        return rabbitTemplate;
    }

    @Bean
    public ConnectionFactory rabbitMQConnectionFactory(){

        Boolean isUriBased = environment.getProperty(URI_BASED_CONNECTION_ENABLED, Boolean.class);
        CachingConnectionFactory connectionFactory;

        if(!Objects.isNull(isUriBased) && isUriBased){
            connectionFactory = new CachingConnectionFactory();
            connectionFactory.setUri(environment.getProperty(RABBITMQ_URI));
        }
        else{
            connectionFactory = new CachingConnectionFactory(rabbitProperties.getHost(), rabbitProperties.getPort());
            connectionFactory.setUsername(rabbitProperties.getUsername());
            connectionFactory.setPassword(rabbitProperties.getPassword());
        }
        return connectionFactory;
    }

    @Bean
    public MessageConverter rabbitMessageConvertor(){
        return new Jackson2JsonMessageConverter();
    }

And publisher code :

public boolean sendMessage(String tenant, T message, int delay){
        MyQueueMessage<T> myQueueMessage = getQueueMessage(tenant, message);
        try{
            amqpTemplate.convertAndSend(exchangeName, routingKey, myQueueMessage, messagePostProcessor -> {

//                        MessageProperties messageProperties = messagePostProcessor.getMessageProperties();
                        messagePostProcessor.getMessageProperties().setHeader("x-message-ttl", 5011);
                        messagePostProcessor.getMessageProperties().setHeader(MessageProperties.X_DELAY, 5012);
                        messagePostProcessor.getMessageProperties().setDelay(5013);
                        messagePostProcessor.getMessageProperties().setReceivedDelay(5014);
                log.info("Setting delay in properties : {}", messagePostProcessor.getMessageProperties().getHeader(MessageProperties.X_DELAY).toString());
                        return messagePostProcessor;
            });
        } catch (Exception e){
            return false;
        }
        return true;
    }

And receiver :

  @RabbitListener(queues = "INVOICE")
    public void receiveMessage(Message message){

    log.info("Message Received : " + message.toString() + " with delay " + message.getMessageProperties().getDelay());
    }
}

Issue : The value

message.getMessageProperties().getDelay() 

comes as NULL in the receiver and the message is also not delayed. It’s getting received instantly.

Did I miss anything?

Please note that I am using docker, rabbitmq-management-3, and have already installed the rabbitmq_delayed_message_exchange plugin.

0 Answers
Related