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 4be02dde27..bfcc065841 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 @@ -63,7 +63,7 @@ import com.rabbitmq.client.Channel; */ public class AmqpInboundGateway extends MessagingGatewaySupport { - private static final ThreadLocal attributesHolder = new ThreadLocal(); + private static final ThreadLocal attributesHolder = new ThreadLocal<>(); private final AbstractMessageListenerContainer messageListenerContainer; @@ -203,11 +203,13 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { @Override protected void doStart() { + super.doStart(); this.messageListenerContainer.start(); } @Override protected void doStop() { + super.doStop(); this.messageListenerContainer.stop(); } @@ -269,18 +271,18 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { org.springframework.messaging.Message converted = convert(message, channel); if (converted != null) { AmqpInboundGateway.this.retryTemplate.execute(context -> { - StaticMessageHeaderAccessor.getDeliveryAttempt(converted).incrementAndGet(); - process(message, converted); - return null; - }, - (RecoveryCallback) AmqpInboundGateway.this.recoveryCallback); + StaticMessageHeaderAccessor.getDeliveryAttempt(converted).incrementAndGet(); + process(message, converted); + return null; + }, + (RecoveryCallback) AmqpInboundGateway.this.recoveryCallback); } } } private org.springframework.messaging.Message convert(Message message, Channel channel) { - Map headers = null; - Object payload = null; + Map headers; + Object payload; boolean isManualAck = AmqpInboundGateway.this.messageListenerContainer .getAcknowledgeMode() == AcknowledgeMode.MANUAL; try { @@ -299,7 +301,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { if (errorChannel != null) { setAttributesIfNecessary(message, null); AmqpInboundGateway.this.messagingTemplate.send(errorChannel, buildErrorMessage(null, - EndpointUtils.errorMessagePayload(message, channel, isManualAck, e))); + EndpointUtils.errorMessagePayload(message, channel, isManualAck, e))); } else { throw e; @@ -307,9 +309,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { return null; } return getMessageBuilderFactory() - .withPayload(payload) - .copyHeaders(headers) - .build(); + .withPayload(payload) + .copyHeaders(headers) + .build(); } private void process(Message message, org.springframework.messaging.Message messagingMessage) { diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java index c2a510ffd1..e7c6fd59f6 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsInboundGateway.java @@ -39,6 +39,7 @@ public class JmsInboundGateway extends MessagingGatewaySupport implements Dispos public JmsInboundGateway(AbstractMessageListenerContainer listenerContainer, ChannelPublishingJmsMessageListener listener) { + this.endpoint = new JmsMessageDrivenEndpoint(listenerContainer, listener); } @@ -138,11 +139,13 @@ public class JmsInboundGateway extends MessagingGatewaySupport implements Dispos @Override protected void doStart() { + super.doStart(); this.endpoint.start(); } @Override protected void doStop() { + super.doStop(); this.endpoint.stop(); }