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 16dd22dfc9..4b164f3191 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 @@ -43,6 +43,7 @@ import org.springframework.integration.amqp.support.EndpointUtils; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.support.ErrorMessageUtils; +import org.springframework.messaging.MessageChannel; import org.springframework.retry.RecoveryCallback; import org.springframework.retry.RetryOperations; import org.springframework.retry.support.RetrySynchronizationManager; @@ -149,8 +150,8 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements /** * Set a batching strategy to use when de-batching messages created by a batching - * producer (such as the BatchingRabbitTemplate). Default is - * {@link SimpleBatchingStrategy}. + * producer (such as the BatchingRabbitTemplate). + * Default is {@link SimpleBatchingStrategy}. * @param batchingStrategy the strategy. * @since 5.2 */ @@ -267,8 +268,8 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements protected final MessageConverter converter = AmqpInboundChannelAdapter.this.messageConverter; // NOSONAR - protected final boolean manualAcks = AcknowledgeMode.MANUAL == - AmqpInboundChannelAdapter.this.messageListenerContainer.getAcknowledgeMode(); // NNOSONAR + protected final boolean manualAcks = // NNOSONAR + AcknowledgeMode.MANUAL == AmqpInboundChannelAdapter.this.messageListenerContainer.getAcknowledgeMode(); protected final RetryOperations retryOps = AmqpInboundChannelAdapter.this.retryTemplate; // NOSONAR @@ -284,7 +285,8 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements createAndSend(message, channel); } else { - final org.springframework.messaging.Message toSend = createMessage(message, channel); + final org.springframework.messaging.Message toSend = + createMessageFromAmqp(message, channel); this.retryOps.execute( context -> { StaticMessageHeaderAccessor.getDeliveryAttempt(toSend).incrementAndGet(); @@ -295,11 +297,13 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements } } catch (MessageConversionException e) { - if (getErrorChannel() != null) { + MessageChannel errorChannel = getErrorChannel(); + if (errorChannel != null) { setAttributesIfNecessary(message, null); getMessagingTemplate() - .send(getErrorChannel(), buildErrorMessage(null, - EndpointUtils.errorMessagePayload(message, channel, this.manualAcks, e))); + .send(errorChannel, + buildErrorMessage(null, + EndpointUtils.errorMessagePayload(message, channel, this.manualAcks, e))); } else { throw e; @@ -313,28 +317,30 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements } private void createAndSend(Message message, Channel channel) { - org.springframework.messaging.Message messagingMessage = createMessage(message, channel); + org.springframework.messaging.Message messagingMessage = createMessageFromAmqp(message, channel); setAttributesIfNecessary(message, messagingMessage); sendMessage(messagingMessage); } - protected org.springframework.messaging.Message createMessage(Message message, Channel channel) { + protected org.springframework.messaging.Message createMessageFromAmqp(Message message, + Channel channel) { + Object payload = convertPayload(message); - Map headers = AmqpInboundChannelAdapter.this.headerMapper - .toHeadersFromRequest(message.getMessageProperties()); + Map headers = + AmqpInboundChannelAdapter.this.headerMapper.toHeadersFromRequest(message.getMessageProperties()); if (AmqpInboundChannelAdapter.this.bindSourceMessage) { headers.put(IntegrationMessageHeaderAccessor.SOURCE_DATA, message); } long deliveryTag = message.getMessageProperties().getDeliveryTag(); - return finalize(channel, payload, headers, deliveryTag); + return createMessageFromPayload(payload, channel, headers, deliveryTag); } protected Object convertPayload(Message message) { Object payload; if (AmqpInboundChannelAdapter.this.batchingStrategy.canDebatch(message.getMessageProperties())) { List payloads = new ArrayList<>(); - AmqpInboundChannelAdapter.this.batchingStrategy.deBatch(message, fragment -> payloads - .add(this.converter.fromMessage(fragment))); + AmqpInboundChannelAdapter.this.batchingStrategy.deBatch(message, + fragment -> payloads.add(this.converter.fromMessage(fragment))); payload = payloads; } else { @@ -343,8 +349,8 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements return payload; } - protected org.springframework.messaging.Message finalize(Channel channel, Object payload, - Map headers, long deliveryTag) { + protected org.springframework.messaging.Message createMessageFromPayload(Object payload, + Channel channel, Map headers, long deliveryTag) { if (this.manualAcks) { headers.put(AmqpHeaders.DELIVERY_TAG, deliveryTag); @@ -375,8 +381,9 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements converted = convertPayloads(messages, channel); } if (converted != null) { - org.springframework.messaging.Message message = finalize(channel, converted, new HashMap<>(), - messages.get(messages.size() - 1).getMessageProperties().getDeliveryTag()); + org.springframework.messaging.Message message = + createMessageFromPayload(converted, channel, new HashMap<>(), + messages.get(messages.size() - 1).getMessageProperties().getDeliveryTag()); try { if (this.retryOps == null) { setAttributesIfNecessary(messages, message); @@ -412,16 +419,15 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements List> converted = new ArrayList<>(); try { - messages.forEach(message -> { - converted.add(createMessage(message, channel)); - }); + messages.forEach(message -> converted.add(createMessageFromAmqp(message, channel))); return converted; } catch (MessageConversionException e) { - if (getErrorChannel() != null) { + MessageChannel errorChannel = getErrorChannel(); + if (errorChannel != null) { setAttributesIfNecessary(messages, null); getMessagingTemplate() - .send(getErrorChannel(), buildErrorMessage(null, + .send(errorChannel, buildErrorMessage(null, EndpointUtils.errorMessagePayload(messages, channel, this.manualAcks, e))); } else { @@ -434,16 +440,15 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements private List convertPayloads(List messages, Channel channel) { List converted = new ArrayList<>(); try { - messages.forEach(message -> { - converted.add(this.converter.fromMessage(message)); - }); + messages.forEach(message -> converted.add(this.converter.fromMessage(message))); return converted; } catch (MessageConversionException e) { - if (getErrorChannel() != null) { + MessageChannel errorChannel = getErrorChannel(); + if (errorChannel != null) { setAttributesIfNecessary(messages, null); getMessagingTemplate() - .send(getErrorChannel(), buildErrorMessage(null, + .send(errorChannel, buildErrorMessage(null, EndpointUtils.errorMessagePayload(messages, channel, this.manualAcks, e))); } else {