From 3e39514003652618967dc0e2658fa0eb98d4f121 Mon Sep 17 00:00:00 2001 From: rahulbats Date: Wed, 7 Nov 2018 12:28:50 -0600 Subject: [PATCH] Add log message for deserialization error handler Add missing logging statement for the logAndContinue deserialization error handler in Kafka Streams binder. Polishing. Resolves #494 --- .../kafka/streams/KafkaStreamsMessageConversionDelegate.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 7417723bf..ea4bfcc4e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -231,10 +231,11 @@ public class KafkaStreamsMessageConversionDelegate { } } else if (kstreamBinderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndFail) { - throw new IllegalStateException("Inbound deserialization failed."); + throw new IllegalStateException("Inbound deserialization failed. Stopping further processing of records."); } else if (kstreamBinderConfigurationProperties.getSerdeError() == KafkaStreamsBinderConfigurationProperties.SerdeError.logAndContinue) { - //quietly pass through. No action needed, this is similar to log and continue. + //quietly passing through. No action needed, this is similar to log and continue. + LOG.error("Inbound deserialization failed. Skipping this record and continuing."); } } }