From 2858954971d75f7331f9068a22c57d9182fd3f72 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 24 Jun 2021 14:48:01 -0400 Subject: [PATCH] GH-3584: Support spring-amqp 2.3.x and 2.4.x # Fix deprecation warning for Reactor's `limitRequest()` --- build.gradle | 2 +- .../inbound/AmqpInboundChannelAdapter.java | 23 +++++++++++++++---- .../amqp/inbound/AmqpInboundGateway.java | 23 ++++++++++++------- .../endpoint/AbstractPollingEndpoint.java | 5 ++-- 4 files changed, 37 insertions(+), 16 deletions(-) diff --git a/build.gradle b/build.gradle index 3b47f9fcf8..9c0c7f3998 100644 --- a/build.gradle +++ b/build.gradle @@ -100,7 +100,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' springDataVersion = project.hasProperty('springDataVersion') ? project.springDataVersion : '2020.0.10' springKafkaVersion = '2.6.9' 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 d6952fe1a4..e21aeb0205 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.support.AmqpHeaders; @@ -100,7 +101,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(); @@ -116,7 +119,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 " + @@ -125,6 +132,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; } @@ -214,7 +224,7 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements else { messageListener = new Listener(); } - this.messageListenerContainer.setMessageListener(messageListener); + this.messageListenerContainer.setupMessageListener(messageListener); this.messageListenerContainer.afterPropertiesSet(); super.onInit(); } @@ -282,8 +292,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 ae8b6091cc..7b865e3ce5 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.support.AmqpHeaders; import org.springframework.amqp.support.converter.MessageConverter; @@ -67,7 +68,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; @@ -96,17 +99,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"); @@ -122,6 +125,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { this.templateMessageConverter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); } setErrorMessageStrategy(new AmqpMessageHeaderErrorMessageStrategy()); + this.abstractListenerContainer = listenerContainer instanceof AbstractMessageListenerContainer + ? (AbstractMessageListenerContainer) listenerContainer + : null; } @@ -241,7 +247,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { + "send an error message when retries are exhausted"); } Listener messageListener = new Listener(); - this.messageListenerContainer.setMessageListener(messageListener); + this.messageListenerContainer.setupMessageListener(messageListener); this.messageListenerContainer.afterPropertiesSet(); if (!this.amqpTemplateExplicitlySet) { ((RabbitTemplate) this.amqpTemplate).afterPropertiesSet(); @@ -336,8 +342,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<>(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java index d516d8c8f0..dbbb0312e4 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java @@ -367,10 +367,11 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement fluxSink.complete(); } }) - .limitRequest( + .take( this.maxMessagesPerPoll < 0 ? Long.MAX_VALUE - : this.maxMessagesPerPoll) + : this.maxMessagesPerPoll, + true) .subscribeOn(Schedulers.fromExecutor(this.taskExecutor)) .doOnComplete(() -> triggerContext