From 8accf70c1e13a2835f133f10bfec96d313558bff Mon Sep 17 00:00:00 2001 From: linyli Date: Mon, 26 Nov 2018 21:02:13 +0800 Subject: [PATCH 1/8] fix kafka-7443:OffsetOutOfRangeException in restoring state store from changelog topic --- .../streams/processor/internals/StoreChangelogReader.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index 34e6e5cdb6f76..ee558b7b2f64c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -107,6 +107,9 @@ public Collection restore(final RestoringTasks active) { needsInitializing.remove(partition); needsRestoring.remove(partition); + final StateRestorer restorer = stateRestorers.get(partition); + restorer.setCheckpointOffset(StateRestorer.NO_CHECKPOINT); + task.reinitializeStateStoresForPartitions(recoverableException.partitions()); } restoreConsumer.seekToBeginning(partitions); From 2aad0a6cffe01f043589ed694e96c5edbbed5cc6 Mon Sep 17 00:00:00 2001 From: linyli Date: Mon, 26 Nov 2018 21:08:33 +0800 Subject: [PATCH 2/8] KAFKA-7443: OffsetOutOfRangeException in restoring state store from changelog topic when start offset of local checkpoint is smaller than that of changelog topic --- .../kafka/streams/processor/internals/StoreChangelogReader.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index ee558b7b2f64c..a79060cb3393b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -107,7 +107,7 @@ public Collection restore(final RestoringTasks active) { needsInitializing.remove(partition); needsRestoring.remove(partition); - final StateRestorer restorer = stateRestorers.get(partition); + final StateRestorer restorer = stateRestorers.get(partition); restorer.setCheckpointOffset(StateRestorer.NO_CHECKPOINT); task.reinitializeStateStoresForPartitions(recoverableException.partitions()); From 9570bc8f484592ec288a6e91c747e48dcd6d3fa0 Mon Sep 17 00:00:00 2001 From: linyli Date: Tue, 27 Nov 2018 10:10:25 +0800 Subject: [PATCH 3/8] remove the tab characters --- .../kafka/streams/processor/internals/StoreChangelogReader.java | 1 - 1 file changed, 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index a79060cb3393b..fdd9d6c303cad 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -109,7 +109,6 @@ public Collection restore(final RestoringTasks active) { final StateRestorer restorer = stateRestorers.get(partition); restorer.setCheckpointOffset(StateRestorer.NO_CHECKPOINT); - task.reinitializeStateStoresForPartitions(recoverableException.partitions()); } restoreConsumer.seekToBeginning(partitions); From 8f6eed132fde9810bd530b6af75e4aea40a34965 Mon Sep 17 00:00:00 2001 From: linyli Date: Thu, 6 Dec 2018 20:32:37 +0800 Subject: [PATCH 4/8] add unit test --- .../internals/StoreChangelogReaderTest.java | 23 +++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java index 34f0a32b88c36..82e29cc09f1cd 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java @@ -157,6 +157,29 @@ public Set partitions() { assertThat(callback.restored.size(), equalTo(messages)); } + @Test + public void shouldRecoverFromOffsetOutOfRangeExceptionAndRestoreFromStart() { + final int message = 10; + setupConsumer(message, topicPartition); + consumer.updateBeginningOffsets(Collections.singletonMap(topicPartition, 5L)); + + StateRestorer stateRestorer = new StateRestorer(topicPartition, restoreListener, 1L, Long.MAX_VALUE, true, + "storeName"); + changelogReader.register(stateRestorer); + + EasyMock.expect(active.restoringTaskFor(topicPartition)).andStubReturn(task); + EasyMock.replay(active, task); + + // first restore call "fails" since OffsetOutOfRangeException but we should not die with an exception + assertEquals(0, changelogReader.restore(active).size()); + // the stateRestorer startoffset is set to NO_CHECKPOINT in the catch block. + assertEquals(StateRestorer.NO_CHECKPOINT, stateRestorer.startingOffset()); + + //restore the active task again + changelogReader.restore(active); + //the restored size should be equal to message length. + assertThat(callback.restored.size(), equalTo(message)); + } @Test public void shouldRestoreMessagesFromCheckpoint() { From 46e56c450bffabdff7c432dcac919f85b116cd00 Mon Sep 17 00:00:00 2001 From: linyli Date: Sat, 8 Dec 2018 18:19:35 +0800 Subject: [PATCH 5/8] update unit test for Kafka-7443 --- .../kafka/clients/consumer/MockConsumer.java | 9 ++++++++ .../internals/StoreChangelogReaderTest.java | 23 ++++++++++++------- 2 files changed, 24 insertions(+), 8 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java index e43c292e6a31e..dbefac41007a0 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java @@ -24,6 +24,8 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.WakeupException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.time.Duration; import java.util.ArrayList; @@ -61,6 +63,7 @@ public class MockConsumer implements Consumer { private KafkaException exception; private AtomicBoolean wakeup; private boolean closed; + private static final Logger log = LoggerFactory.getLogger(MockConsumer.class); public MockConsumer(OffsetResetStrategy offsetResetStrategy) { this.subscriptions = new SubscriptionState(offsetResetStrategy); @@ -188,6 +191,12 @@ public synchronized ConsumerRecords poll(final Duration timeout) { if (!subscriptions.isPaused(entry.getKey())) { final List> recs = entry.getValue(); for (final ConsumerRecord rec : recs) { + if(rec.offset() > subscriptions.position(entry.getKey())) { + log.info("poll offset{} less than changelog topic start offset {}, throw OffsetOutOfRangeException!", rec.offset(), subscriptions.position(entry.getKey())); + RuntimeException exception = new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), rec.offset())); + throw exception; + } + if (assignment().contains(entry.getKey()) && rec.offset() >= subscriptions.position(entry.getKey())) { results.computeIfAbsent(entry.getKey(), partition -> new ArrayList<>()).add(rec); subscriptions.position(entry.getKey(), rec.offset() + 1); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java index 82e29cc09f1cd..abd41364db4d6 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java @@ -159,11 +159,16 @@ public Set partitions() { @Test public void shouldRecoverFromOffsetOutOfRangeExceptionAndRestoreFromStart() { - final int message = 10; - setupConsumer(message, topicPartition); - consumer.updateBeginningOffsets(Collections.singletonMap(topicPartition, 5L)); + final int messages = 10; + final int startOffset = 5; + assignPartition(messages, topicPartition); + consumer.updateBeginningOffsets(Collections.singletonMap(topicPartition, (long) startOffset)); + consumer.updateEndOffsets(Collections.singletonMap(topicPartition, (long) (messages + startOffset))); + + addRecords(messages, topicPartition, startOffset); + consumer.assign(Collections.emptyList()); - StateRestorer stateRestorer = new StateRestorer(topicPartition, restoreListener, 1L, Long.MAX_VALUE, true, + final StateRestorer stateRestorer = new StateRestorer(topicPartition, restoreListener, 1L, Long.MAX_VALUE, true, "storeName"); changelogReader.register(stateRestorer); @@ -172,13 +177,15 @@ public void shouldRecoverFromOffsetOutOfRangeExceptionAndRestoreFromStart() { // first restore call "fails" since OffsetOutOfRangeException but we should not die with an exception assertEquals(0, changelogReader.restore(active).size()); - // the stateRestorer startoffset is set to NO_CHECKPOINT in the catch block. - assertEquals(StateRestorer.NO_CHECKPOINT, stateRestorer.startingOffset()); + //the starting offset for stateRestorer is set to NO_CHECKPOINT + assertThat(stateRestorer.checkpoint(), equalTo(-1L)); //restore the active task again - changelogReader.restore(active); + changelogReader.register(stateRestorer); + //the restored task should return completed partition without Exception. + assertEquals(1, changelogReader.restore(active).size()); //the restored size should be equal to message length. - assertThat(callback.restored.size(), equalTo(message)); + assertThat(callback.restored.size(), equalTo(messages)); } @Test From d6d58a76578be2fd1795591543901519865c9065 Mon Sep 17 00:00:00 2001 From: linyli Date: Sat, 8 Dec 2018 18:25:40 +0800 Subject: [PATCH 6/8] update exception in MockConsumer.java --- .../java/org/apache/kafka/clients/consumer/MockConsumer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java index dbefac41007a0..8ed73aa9219c1 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java @@ -193,7 +193,7 @@ public synchronized ConsumerRecords poll(final Duration timeout) { for (final ConsumerRecord rec : recs) { if(rec.offset() > subscriptions.position(entry.getKey())) { log.info("poll offset{} less than changelog topic start offset {}, throw OffsetOutOfRangeException!", rec.offset(), subscriptions.position(entry.getKey())); - RuntimeException exception = new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), rec.offset())); + RuntimeException exception = new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), subscriptions.position(entry.getKey()))); throw exception; } From c355325a89f9afe3daf0976385c9032dac019311 Mon Sep 17 00:00:00 2001 From: linyli Date: Sun, 9 Dec 2018 00:07:05 +0800 Subject: [PATCH 7/8] fix test failed --- .../java/org/apache/kafka/clients/consumer/MockConsumer.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java index 8ed73aa9219c1..fe142900909ce 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java @@ -191,8 +191,8 @@ public synchronized ConsumerRecords poll(final Duration timeout) { if (!subscriptions.isPaused(entry.getKey())) { final List> recs = entry.getValue(); for (final ConsumerRecord rec : recs) { - if(rec.offset() > subscriptions.position(entry.getKey())) { - log.info("poll offset{} less than changelog topic start offset {}, throw OffsetOutOfRangeException!", rec.offset(), subscriptions.position(entry.getKey())); + if (beginningOffsets.get(entry.getKey()) != null && beginningOffsets.get(entry.getKey()) > subscriptions.position(entry.getKey())) { + log.info("poll offset {} less than changelog topic start offset {}, throw OffsetOutOfRangeException!", subscriptions.position(entry.getKey()), beginningOffsets.get(entry.getKey())); RuntimeException exception = new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), subscriptions.position(entry.getKey()))); throw exception; } From 18fc9f98f90870a52c491b5c3b22c98f4fd8276a Mon Sep 17 00:00:00 2001 From: linyli Date: Mon, 10 Dec 2018 11:37:55 +0800 Subject: [PATCH 8/8] modify exception and some text format --- .../org/apache/kafka/clients/consumer/MockConsumer.java | 7 +------ .../processor/internals/StoreChangelogReaderTest.java | 8 +++++++- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java index fe142900909ce..f877f9d13a451 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java @@ -24,8 +24,6 @@ import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.WakeupException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import java.time.Duration; import java.util.ArrayList; @@ -63,7 +61,6 @@ public class MockConsumer implements Consumer { private KafkaException exception; private AtomicBoolean wakeup; private boolean closed; - private static final Logger log = LoggerFactory.getLogger(MockConsumer.class); public MockConsumer(OffsetResetStrategy offsetResetStrategy) { this.subscriptions = new SubscriptionState(offsetResetStrategy); @@ -192,9 +189,7 @@ public synchronized ConsumerRecords poll(final Duration timeout) { final List> recs = entry.getValue(); for (final ConsumerRecord rec : recs) { if (beginningOffsets.get(entry.getKey()) != null && beginningOffsets.get(entry.getKey()) > subscriptions.position(entry.getKey())) { - log.info("poll offset {} less than changelog topic start offset {}, throw OffsetOutOfRangeException!", subscriptions.position(entry.getKey()), beginningOffsets.get(entry.getKey())); - RuntimeException exception = new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), subscriptions.position(entry.getKey()))); - throw exception; + throw new OffsetOutOfRangeException(Collections.singletonMap(entry.getKey(), subscriptions.position(entry.getKey()))); } if (assignment().contains(entry.getKey()) && rec.offset() >= subscriptions.position(entry.getKey())) { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java index abd41364db4d6..d08f0d7360d18 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java @@ -161,6 +161,7 @@ public Set partitions() { public void shouldRecoverFromOffsetOutOfRangeExceptionAndRestoreFromStart() { final int messages = 10; final int startOffset = 5; + final long expiredCheckpoint = 1L; assignPartition(messages, topicPartition); consumer.updateBeginningOffsets(Collections.singletonMap(topicPartition, (long) startOffset)); consumer.updateEndOffsets(Collections.singletonMap(topicPartition, (long) (messages + startOffset))); @@ -168,7 +169,12 @@ public void shouldRecoverFromOffsetOutOfRangeExceptionAndRestoreFromStart() { addRecords(messages, topicPartition, startOffset); consumer.assign(Collections.emptyList()); - final StateRestorer stateRestorer = new StateRestorer(topicPartition, restoreListener, 1L, Long.MAX_VALUE, true, + final StateRestorer stateRestorer = new StateRestorer( + topicPartition, + restoreListener, + expiredCheckpoint, + Long.MAX_VALUE, + true, "storeName"); changelogReader.register(stateRestorer);