diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapter.java index 6667a7b9..3bab4dd7 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapter.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -25,6 +25,7 @@ import org.springframework.amqp.rabbit.batch.SimpleBatchingStrategy; import org.springframework.amqp.rabbit.listener.api.ChannelAwareBatchMessageListener; import org.springframework.amqp.rabbit.listener.api.RabbitListenerErrorHandler; import org.springframework.amqp.rabbit.support.RabbitExceptionTranslator; +import org.springframework.amqp.support.converter.MessageConversionException; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; @@ -63,7 +64,20 @@ public class BatchMessagingMessageListenerAdapter extends MessagingMessageListen else { List> messagingMessages = new ArrayList<>(); for (org.springframework.amqp.core.Message message : messages) { - messagingMessages.add(toMessagingMessage(message)); + try { + Message messagingMessage = toMessagingMessage(message); + messagingMessages.add(messagingMessage); + } + catch (MessageConversionException e) { + this.logger.error("Could not convert incoming message", e); + try { + channel.basicReject(message.getMessageProperties().getDeliveryTag(), false); + } + catch (Exception ex) { + this.logger.error("Failed to reject message with conversion error", ex); + throw e; // NOSONAR + } + } } if (this.converterAdapter.isMessageList()) { converted = new GenericMessage<>(messagingMessages); diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapterTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapterTests.java index 358c4676..d12f1d5a 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapterTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/adapter/BatchMessagingMessageListenerAdapterTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2022 the original author or authors. + * Copyright 2022-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -21,16 +21,42 @@ import static org.assertj.core.api.Assertions.assertThatIllegalStateException; import java.lang.reflect.Method; import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import org.junit.jupiter.api.Test; +import org.springframework.amqp.core.Message; +import org.springframework.amqp.core.MessageBuilder; +import org.springframework.amqp.core.MessagePropertiesBuilder; +import org.springframework.amqp.rabbit.annotation.EnableRabbit; +import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory; +import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.amqp.rabbit.junit.RabbitAvailable; +import org.springframework.amqp.rabbit.junit.RabbitAvailableCondition; +import org.springframework.amqp.rabbit.listener.RabbitListenerContainerFactory; +import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; +import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.utils.test.TestUtils; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; + +import com.fasterxml.jackson.databind.ObjectMapper; /** * @author Gary Russell + * @author heng zhang + * * @since 3.0 * */ +@SpringJUnitConfig +@RabbitAvailable(queues = "test.batchQueue") public class BatchMessagingMessageListenerAdapterTests { @Test @@ -52,4 +78,104 @@ public class BatchMessagingMessageListenerAdapterTests { public void listen(List in) { } + + @Test + public void errorMsgConvert(@Autowired BatchMessagingMessageListenerAdapterTests.Config config, + @Autowired RabbitTemplate template) throws Exception { + + Message message = MessageBuilder.withBody(""" + { + "name" : "Tom", + "age" : 18 + } + """.getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build(); + + Message errorMessage = MessageBuilder.withBody("".getBytes()).andProperties( + MessagePropertiesBuilder.newInstance() + .setContentType("application/json") + .setReplyTo("nowhere") + .build()) + .build(); + + for (int i = 0; i < config.count; i++) { + template.send("test.batchQueue", message); + template.send("test.batchQueue", errorMessage); + } + + assertThat(config.countDownLatch.await(config.count * 1000L, TimeUnit.SECONDS)).isTrue(); + } + + + + @Configuration + @EnableRabbit + public static class Config { + volatile int count = 5; + volatile CountDownLatch countDownLatch = new CountDownLatch(count); + + @RabbitListener( + queues = "test.batchQueue", + containerFactory = "batchListenerContainerFactory" + ) + public void listen(List list) { + for (Model model : list) { + countDownLatch.countDown(); + } + + } + + @Bean + ConnectionFactory cf() { + return new CachingConnectionFactory(RabbitAvailableCondition.getBrokerRunning().getConnectionFactory()); + } + + @Bean(name = "batchListenerContainerFactory") + public RabbitListenerContainerFactory rc(ConnectionFactory connectionFactory) { + SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); + factory.setConnectionFactory(connectionFactory); + factory.setPrefetchCount(1); + factory.setConcurrentConsumers(1); + factory.setBatchListener(true); + factory.setBatchSize(3); + factory.setConsumerBatchEnabled(true); + + Jackson2JsonMessageConverter jackson2JsonMessageConverter = new Jackson2JsonMessageConverter(new ObjectMapper()); + factory.setMessageConverter(jackson2JsonMessageConverter); + + return factory; + } + + @Bean + RabbitTemplate template(ConnectionFactory cf) { + return new RabbitTemplate(cf); + } + + + } + public static class Model { + String name; + String age; + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + + public String getAge() { + return age; + } + + public void setAge(String age) { + this.age = age; + } + } + }