diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java index 8523edfc..3587793f 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/RabbitTemplate.java @@ -2631,6 +2631,7 @@ public class RabbitTemplate extends RabbitAccessor // NOSONAR type line count future.completeExceptionally( new ConsumeOkNotReceivedException("Blocking receive, consumer failed to consume within " + timeoutMillis + " ms: " + consumer)); + RabbitUtils.setPhysicalCloseRequired(channel, true); } return consumer; } diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java index d9d82c90..5b05a99f 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/AbstractMessageListenerContainer.java @@ -42,6 +42,7 @@ import org.springframework.amqp.AmqpRejectAndDontRequeueException; import org.springframework.amqp.ImmediateAcknowledgeAmqpException; import org.springframework.amqp.core.AcknowledgeMode; import org.springframework.amqp.core.AmqpAdmin; +import org.springframework.amqp.core.BatchMessageListener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageListener; import org.springframework.amqp.core.MessagePostProcessor; @@ -1920,6 +1921,17 @@ public abstract class AbstractMessageListenerContainer extends RabbitAccessor } } + @Nullable + protected List debatch(Message message) { + if (isDeBatchingEnabled() && getBatchingStrategy().canDebatch(message.getMessageProperties()) + && getMessageListener() instanceof BatchMessageListener) { + final List messageList = new ArrayList<>(); + getBatchingStrategy().deBatch(message, fragment -> messageList.add(fragment)); + return messageList; + } + return null; + } + @FunctionalInterface private interface ContainerDelegate { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java index 62be2924..271f4f99 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/DirectMessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2019 the original author or authors. + * Copyright 2016-2020 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. @@ -974,9 +974,14 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta this.logger.debug(this + " received " + message); } updateLastReceive(); + Object data = message; + List debatched = debatch(message); + if (debatched != null) { + data = debatched; + } if (this.transactionManager != null) { try { - executeListenerInTransaction(message, deliveryTag); + executeListenerInTransaction(data, deliveryTag); } catch (WrappedTransactionException e) { if (e.getCause() instanceof Error) { @@ -994,7 +999,7 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } else { try { - callExecuteListener(message, deliveryTag); + callExecuteListener(data, deliveryTag); } catch (Exception e) { // NOSONAR @@ -1002,7 +1007,7 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } } - private void executeListenerInTransaction(Message message, long deliveryTag) { + private void executeListenerInTransaction(Object data, long deliveryTag) { if (this.isRabbitTxManager) { ConsumerChannelRegistry.registerConsumerChannel(getChannel(), this.connectionFactory); } @@ -1018,7 +1023,7 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta } // unbound in ResourceHolderSynchronization.beforeCompletion() try { - callExecuteListener(message, deliveryTag); + callExecuteListener(data, deliveryTag); } catch (RuntimeException e1) { prepareHolderForRollback(resourceHolder, e1); @@ -1031,10 +1036,10 @@ public class DirectMessageListenerContainer extends AbstractMessageListenerConta }); } - private void callExecuteListener(Message message, long deliveryTag) { + private void callExecuteListener(Object data, long deliveryTag) { boolean channelLocallyTransacted = isChannelLocallyTransacted(); try { - executeListener(getChannel(), message); + executeListener(getChannel(), data); handleAck(deliveryTag, channelLocallyTransacted); } catch (ImmediateAcknowledgeAmqpException e) { diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java index 781a2e30..ff333848 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java @@ -945,6 +945,10 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } } else { + messages = debatch(message); + if (messages != null) { + break; + } try { executeListener(channel, message); } @@ -994,7 +998,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } } } - if (this.consumerBatchEnabled && messages != null) { + if (messages != null) { executeWithList(channel, messages, deliveryTag, consumer); } diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplateTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplateTests.java index bb4eaaca..d36b5152 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplateTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplateTests.java @@ -42,6 +42,7 @@ import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.springframework.amqp.AmqpException; +import org.springframework.amqp.core.BatchMessageListener; import org.springframework.amqp.core.Message; import org.springframework.amqp.core.MessageDeliveryMode; import org.springframework.amqp.core.MessageListener; @@ -53,7 +54,9 @@ import org.springframework.amqp.rabbit.connection.ThreadChannelConnectionFactory import org.springframework.amqp.rabbit.junit.BrokerTestUtils; import org.springframework.amqp.rabbit.junit.RabbitAvailable; import org.springframework.amqp.rabbit.junit.RabbitAvailableCondition; +import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; import org.springframework.amqp.rabbit.listener.ConditionalRejectingErrorHandler; +import org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.support.postprocessor.AbstractCompressingPostProcessor; @@ -228,20 +231,47 @@ public class BatchingRabbitTemplateTests { } @Test - public void testDebatchByContainer() throws Exception { - final List received = new ArrayList(); - final CountDownLatch latch = new CountDownLatch(2); + void testDebatchSMLCSplit() throws Exception { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.connectionFactory); + container.setReceiveTimeout(100); + testDebatchByContainer(container, false); + } + + @Test + void testDebatchSMLC() throws Exception { + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.connectionFactory); + container.setReceiveTimeout(100); + testDebatchByContainer(container, true); + } + + @Test + void testDebatchDMLC() throws Exception { + testDebatchByContainer(new DirectMessageListenerContainer(this.connectionFactory), true); + } + + private void testDebatchByContainer(AbstractMessageListenerContainer container, boolean asList) throws Exception { + final List received = new ArrayList(); + final CountDownLatch latch = new CountDownLatch(asList ? 1 : 2); container.setQueueNames(ROUTE); List lastInBatch = new ArrayList<>(); AtomicInteger batchSize = new AtomicInteger(); - container.setMessageListener((MessageListener) message -> { - received.add(message); - lastInBatch.add(message.getMessageProperties().isLastInBatch()); - batchSize.set(message.getMessageProperties().getHeader(AmqpHeaders.BATCH_SIZE)); - latch.countDown(); - }); - container.setReceiveTimeout(100); + if (asList) { + container.setMessageListener((BatchMessageListener) messages -> { + received.addAll(messages); + lastInBatch.add(false); + lastInBatch.add(true); + batchSize.set(messages.size()); + latch.countDown(); + }); + } + else { + container.setMessageListener((MessageListener) message -> { + received.add(message); + lastInBatch.add(message.getMessageProperties().isLastInBatch()); + batchSize.set(message.getMessageProperties().getHeader(AmqpHeaders.BATCH_SIZE)); + latch.countDown(); + }); + } container.afterPropertiesSet(); container.start(); try { diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 966e4c53..0eee4c73 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -2001,11 +2001,12 @@ Batched messages (created by a producer) are automatically de-batched by listene Rejecting any message from a batch causes the entire batch to be rejected. See <> for more information about batching. -Starting with version 2.2, the `SimpleMessageListeneContainer` can be use to create batches on the consumer side (where the producer sent discrete messages). +Starting with version 2.2, the `SimpleMessageListenerContainer` can be use to create batches on the consumer side (where the producer sent discrete messages). Set the container property `consumerBatchEnabled` to enable this feature. `deBatchingEnabled` must also be true so that the container is responsible for processing batches of both types. Implement `BatchMessageListener` or `ChannelAwareBatchMessageListener` when `consumerBatchEnabled` is true. +Starting with version 2.2.7 both the `SimpleMessageListenerContainer` and `DirectMessageListenerContainer` can debatch <> as `List`. See <> for information about using this feature with `@RabbitListener`. [[consumer-events]] @@ -5472,6 +5473,8 @@ a|image::images/tickmark.png[] (N/A) |When true, the listener container will debatch batched messages and invoke the listener with each message from the batch. +Starting with version 2.2.7, <> will be debatched as a `List` if the listener is a `BatchMessageListener` or `ChannelAwareBatchMessageListener`. +Otherwise messages from the batch are presented one-at-a-time. Default true. See <> and <>.