GH-3584: Support spring-amqp 2.3.x and 2.4.x

# Fix deprecation warning for Reactor's `limitRequest()`
This commit is contained in:
Gary Russell
2021-06-24 14:48:01 -04:00
committed by Artem Bilan
parent ddab28a06f
commit 2858954971
4 changed files with 37 additions and 16 deletions

View File

@@ -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'

View File

@@ -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<AttributeAccessor> 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

View File

@@ -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<AttributeAccessor> 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<Object> convert(Message message, Channel channel) {
Map<String, Object> 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<Object> payloads = new ArrayList<>();

View File

@@ -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