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 dcad7140..15719669 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 @@ -41,7 +41,7 @@ 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 " @@ -50,7 +50,7 @@ public class SeekToCurrentErrorHandler implements ContainerAwareErrorHandler { Map offsets = new LinkedHashMap<>(); records.forEach(r -> - offsets.computeIfAbsent(new TopicPartition(r.topic(), r.partition()), k -> r.offset())); + offsets.computeIfAbsent(new TopicPartition(r.topic(), r.partition()), k -> r.offset())); offsets.forEach(consumer::seek); throw new KafkaException("Seek to current after exception", thrownException); }