From adb2f9a2314fd48196e006f439708403edb2f6e2 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 29 Oct 2019 16:31:45 -0400 Subject: [PATCH] GH-1287: Detect mis-configured deserialization Resolves https://github.com/spring-projects/spring-kafka/issues/1287 Throw an `IllegalStateException` in the `SeekToCurrentErrorHandler` if the root exception is a `SerializationException`. Handling of deserializatin problms requires the configuration of an `ErrorHandlingDeserializer2`. **cherry-pick to 2.2.x, 2.1.x** # Conflicts: # spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java # spring-kafka/src/test/java/org/springframework/kafka/listener/SeekToCurrentErrorHandlerTests.java # Conflicts: # spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java --- .../kafka/listener/SeekToCurrentErrorHandler.java | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java index 9de509e1..dcad7140 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/SeekToCurrentErrorHandler.java @@ -23,6 +23,7 @@ import java.util.Map; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.SerializationException; import org.springframework.kafka.KafkaException; @@ -40,6 +41,13 @@ public class SeekToCurrentErrorHandler implements ContainerAwareErrorHandler { @Override public void handle(Exception thrownException, List> records, Consumer consumer, MessageListenerContainer container) { + + if (thrownException instanceof SerializationException) { + throw new IllegalStateException("This error handler cannot process 'SerializationException's directly, " + + "please consider configuring an 'ErrorHandlingDeserializer2' in the value and/or key " + + "deserializer", thrownException); + } + Map offsets = new LinkedHashMap<>(); records.forEach(r -> offsets.computeIfAbsent(new TopicPartition(r.topic(), r.partition()), k -> r.offset()));