From e16ae4fc0a53622acc5117bdb49e21dc5e721d7a Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 28 Aug 2018 12:02:23 -0400 Subject: [PATCH] SeekToCurrentErrorHandler and Recovery Don't throw an exception if the record was skipped by invoking the recoverer. - Adds noise to the log since the record was recovered in some way. --- .../listener/SeekToCurrentErrorHandler.java | 5 ++-- .../kafka/support/SeekUtils.java | 20 +++++++++++--- .../listener/SeekToCurrentRecovererTests.java | 27 +++++++++++++++++++ 3 files changed, 47 insertions(+), 5 deletions(-) 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 883bbd37..5e89e10a 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 @@ -78,8 +78,9 @@ public class SeekToCurrentErrorHandler implements ContainerAwareErrorHandler { @Override public void handle(Exception thrownException, List> records, Consumer consumer, MessageListenerContainer container) { - SeekUtils.doSeeks(records, consumer, thrownException, true, this.failureTracker::skip, logger); - throw new KafkaException("Seek to current after exception", thrownException); + if (!SeekUtils.doSeeks(records, consumer, thrownException, true, this.failureTracker::skip, logger)) { + throw new KafkaException("Seek to current after exception", thrownException); + } } @Override diff --git a/spring-kafka/src/main/java/org/springframework/kafka/support/SeekUtils.java b/spring-kafka/src/main/java/org/springframework/kafka/support/SeekUtils.java index 3d1cd6c2..de015a7b 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/support/SeekUtils.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/support/SeekUtils.java @@ -53,25 +53,39 @@ public final class SeekUtils { * @param recoverable true if skipping the first record is allowed. * @param skipper function to determine whether or not to skip seeking the first. * @param logger a {@link Log} for seek errors. + * @return true if the failed record was skipped. */ - public static void doSeeks(List> records, Consumer consumer, Exception exception, + public static boolean doSeeks(List> records, Consumer consumer, Exception exception, boolean recoverable, BiPredicate, Exception> skipper, Log logger) { Map partitions = new LinkedHashMap<>(); AtomicBoolean first = new AtomicBoolean(true); + AtomicBoolean skipped = new AtomicBoolean(); records.forEach(record -> { - if (!recoverable || !first.get() || !skipper.test(record, exception)) { - partitions.computeIfAbsent(new TopicPartition(record.topic(), record.partition()), offset -> record.offset()); + if (recoverable && first.get()) { + skipped.set(skipper.test(record, exception)); + if (skipped.get() && logger.isDebugEnabled()) { + logger.debug("Skipping seek of: " + record); + } + } + if (!recoverable || !first.get() || !skipped.get()) { + partitions.computeIfAbsent(new TopicPartition(record.topic(), record.partition()), + offset -> record.offset()); } first.set(false); }); + boolean tracing = logger.isTraceEnabled(); partitions.forEach((topicPartition, offset) -> { try { + if (tracing) { + logger.trace("Seeking: " + topicPartition + " to: " + offset); + } consumer.seek(topicPartition, offset); } catch (Exception e) { logger.error("Failed to seek " + topicPartition + " to " + offset, e); } }); + return skipped.get(); } } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/SeekToCurrentRecovererTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/SeekToCurrentRecovererTests.java index 3b2fde2c..e24202b9 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/SeekToCurrentRecovererTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/SeekToCurrentRecovererTests.java @@ -17,9 +17,14 @@ package org.springframework.kafka.listener; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.fail; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; +import java.util.ArrayList; +import java.util.List; import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -33,6 +38,7 @@ import org.apache.kafka.common.TopicPartition; import org.junit.ClassRule; import org.junit.Test; +import org.springframework.kafka.KafkaException; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -121,4 +127,25 @@ public class SeekToCurrentRecovererTests { verify(errorHandler).clearThreadState(); } + @Test + public void seekToCurrentErrorHandlerRecovers() { + SeekToCurrentErrorHandler eh = new SeekToCurrentErrorHandler((r, e) -> { }, 2); + List> records = new ArrayList<>(); + records.add(new ConsumerRecord<>("foo", 0, 0, null, "foo")); + records.add(new ConsumerRecord<>("foo", 0, 1, null, "bar")); + Consumer consumer = mock(Consumer.class); + try { + eh.handle(new RuntimeException(), records, consumer, null); + fail("Expected exception"); + } + catch (KafkaException e) { + // NOSONAR + } + verify(consumer).seek(new TopicPartition("foo", 0), 0L); + verifyNoMoreInteractions(consumer); + eh.handle(new RuntimeException(), records, consumer, null); + verify(consumer).seek(new TopicPartition("foo", 0), 1L); + verifyNoMoreInteractions(consumer); + } + }