diff --git a/build.gradle b/build.gradle index 027e32ed5c..aabcff70fe 100644 --- a/build.gradle +++ b/build.gradle @@ -98,7 +98,7 @@ ext { servletApiVersion = '4.0.1' smackVersion = '4.3.5' soapVersion = '1.4.0' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.3.9' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.3.10-SNAPSHOT' springDataVersion = project.hasProperty('springDataVersion') ? project.springDataVersion : '2021.0.2' springKafkaVersion = '2.7.3' springRetryVersion = '1.3.1' diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java index 00f41ff7e7..0788f0d1c5 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java @@ -27,6 +27,7 @@ import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.batch.BatchingStrategy; import org.springframework.amqp.rabbit.batch.SimpleBatchingStrategy; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.MessageListenerContainer; import org.springframework.amqp.rabbit.listener.api.ChannelAwareBatchMessageListener; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.retry.MessageBatchRecoverer; @@ -102,7 +103,9 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements private static final ThreadLocal ATTRIBUTES_HOLDER = new ThreadLocal<>(); - private final AbstractMessageListenerContainer messageListenerContainer; + private final MessageListenerContainer messageListenerContainer; + + private final AbstractMessageListenerContainer abstractListenerContainer; private MessageConverter messageConverter = new SimpleMessageConverter(); @@ -120,7 +123,11 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements private BatchMode batchMode = BatchMode.MESSAGES; - public AmqpInboundChannelAdapter(AbstractMessageListenerContainer listenerContainer) { + /** + * Construct an instance using the provided container. + * @param listenerContainer the container. + */ + public AmqpInboundChannelAdapter(MessageListenerContainer listenerContainer) { Assert.notNull(listenerContainer, "listenerContainer must not be null"); Assert.isNull(listenerContainer.getMessageListener(), "The listenerContainer provided to an AMQP inbound Channel Adapter " + @@ -129,6 +136,9 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements this.messageListenerContainer = listenerContainer; this.messageListenerContainer.setAutoStartup(false); setErrorMessageStrategy(new AmqpMessageHeaderErrorMessageStrategy()); + this.abstractListenerContainer = listenerContainer instanceof AbstractMessageListenerContainer + ? (AbstractMessageListenerContainer) listenerContainer + : null; } @@ -230,7 +240,7 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements else { messageListener = new Listener(); } - this.messageListenerContainer.setMessageListener(messageListener); + this.messageListenerContainer.setupMessageListener(messageListener); this.messageListenerContainer.afterPropertiesSet(); super.onInit(); } @@ -330,8 +340,11 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements protected final MessageConverter converter = AmqpInboundChannelAdapter.this.messageConverter; // NOSONAR - protected final boolean manualAcks = // NNOSONAR - AcknowledgeMode.MANUAL == AmqpInboundChannelAdapter.this.messageListenerContainer.getAcknowledgeMode(); + protected final boolean manualAcks = // NOSONAR + AmqpInboundChannelAdapter.this.abstractListenerContainer == null + ? false + : AcknowledgeMode.MANUAL == AmqpInboundChannelAdapter.this.abstractListenerContainer + .getAcknowledgeMode(); protected final RetryOperations retryOps = AmqpInboundChannelAdapter.this.retryTemplate; // NOSONAR diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java index 2a34203410..7ab573923c 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java @@ -29,6 +29,7 @@ import org.springframework.amqp.rabbit.batch.BatchingStrategy; import org.springframework.amqp.rabbit.batch.SimpleBatchingStrategy; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; +import org.springframework.amqp.rabbit.listener.MessageListenerContainer; import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.retry.MessageRecoverer; import org.springframework.amqp.support.AmqpHeaders; @@ -68,7 +69,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { private static final ThreadLocal ATTRIBUTES_HOLDER = new ThreadLocal<>(); - private final AbstractMessageListenerContainer messageListenerContainer; + private final MessageListenerContainer messageListenerContainer; + + private final AbstractMessageListenerContainer abstractListenerContainer; private final AmqpTemplate amqpTemplate; @@ -99,17 +102,17 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { } /** - * Construct {@link AmqpInboundGateway} based on the provided {@link AbstractMessageListenerContainer} + * Construct {@link AmqpInboundGateway} based on the provided {@link MessageListenerContainer} * to receive request messages and {@link AmqpTemplate} to send replies. - * @param listenerContainer the {@link AbstractMessageListenerContainer} to receive AMQP messages. + * @param listenerContainer the {@link MessageListenerContainer} to receive AMQP messages. * @param amqpTemplate the {@link AmqpTemplate} to send reply messages. * @since 4.2 */ - public AmqpInboundGateway(AbstractMessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate) { + public AmqpInboundGateway(MessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate) { this(listenerContainer, amqpTemplate, true); } - private AmqpInboundGateway(AbstractMessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate, + private AmqpInboundGateway(MessageListenerContainer listenerContainer, AmqpTemplate amqpTemplate, boolean amqpTemplateExplicitlySet) { Assert.notNull(listenerContainer, "listenerContainer must not be null"); Assert.notNull(amqpTemplate, "'amqpTemplate' must not be null"); @@ -125,6 +128,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { this.templateMessageConverter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); } setErrorMessageStrategy(new AmqpMessageHeaderErrorMessageStrategy()); + this.abstractListenerContainer = listenerContainer instanceof AbstractMessageListenerContainer + ? (AbstractMessageListenerContainer) listenerContainer + : null; } @@ -256,7 +262,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { setupRecoveryCallbackIfAny(); } Listener messageListener = new Listener(); - this.messageListenerContainer.setMessageListener(messageListener); + this.messageListenerContainer.setupMessageListener(messageListener); this.messageListenerContainer.afterPropertiesSet(); if (!this.amqpTemplateExplicitlySet) { ((RabbitTemplate) this.amqpTemplate).afterPropertiesSet(); @@ -366,8 +372,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { private org.springframework.messaging.Message convert(Message message, Channel channel) { Map headers; Object payload; - boolean isManualAck = - AmqpInboundGateway.this.messageListenerContainer.getAcknowledgeMode() == AcknowledgeMode.MANUAL; + boolean isManualAck = AmqpInboundGateway.this.abstractListenerContainer == null + ? false + : AcknowledgeMode.MANUAL == AmqpInboundGateway.this.abstractListenerContainer.getAcknowledgeMode(); try { if (AmqpInboundGateway.this.batchingStrategy.canDebatch(message.getMessageProperties())) { List payloads = new ArrayList<>();