diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index 66b3c9968f..818b0a7eb0 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -443,8 +443,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo @Override public void onMessage(ConsumerRecord record, Acknowledgment acknowledgment, Consumer consumer) { - - Message message = null; + Message message; try { message = toMessagingMessage(record, acknowledgment, consumer); } @@ -452,22 +451,22 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo if (KafkaMessageDrivenChannelAdapter.this.retryTemplate == null) { setAttributesIfNecessary(record, null, true); } - MessageChannel errorChannel = getErrorChannel(); - if (errorChannel != null) { - RuntimeException exception = new ConversionException("Failed to convert to message", record, ex); - sendErrorMessageIfNecessary(null, exception); + + RuntimeException exception = new ConversionException("Failed to convert to message", record, ex); + if (sendErrorMessageIfNecessary(null, exception)) { + return; } else { throw ex; } } + RetryTemplate template = KafkaMessageDrivenChannelAdapter.this.retryTemplate; if (template != null) { - Message toSend = message; doWithRetry(template, KafkaMessageDrivenChannelAdapter.this.recoveryCallback, record, acknowledgment, consumer, () -> { if (!KafkaMessageDrivenChannelAdapter.this.filterInRetry || passesFilter(record)) { - sendMessageIfAny(enhanceHeadersAndSaveAttributes(toSend, record), record); + sendMessageIfAny(enhanceHeadersAndSaveAttributes(message, record), record); } }); }