GH-3584: Support spring-amqp 2.3.x and 2.4.x
This commit is contained in:
@@ -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'
|
||||
|
||||
@@ -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<AttributeAccessor> 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
|
||||
|
||||
|
||||
@@ -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<AttributeAccessor> 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<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<>();
|
||||
|
||||
Reference in New Issue
Block a user