From cdcc986d1148e5dbecd80b073c76b483e725bc74 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 31 May 2022 14:08:07 -0400 Subject: [PATCH] Fix KafkaMDChannelAdapter error handling Current `IntegrationRecordMessageListener.onMessage()` send a conversion error to the `errorChannel`, but it does not return if it was successful leading to the NPE `enhanceHeadersAndSaveAttributes()` because the `message` is null * Change the `IntegrationRecordMessageListener.onMessage()` logic to check the result of the `sendErrorMessageIfNecessary()` and return immediately if success. Rethrow an exception otherwise. --- .../inbound/KafkaMessageDrivenChannelAdapter.java | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) 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); } }); }