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."); } } }