From eb859861e904f982e987f151257f220608e66dca Mon Sep 17 00:00:00 2001 From: pgwhalen Date: Sun, 17 Feb 2019 19:12:26 -0600 Subject: [PATCH 1/5] KAFKA-7941: Catch TimeoutException in KafkaBasedLog worker thread - When calling readLogToEnd(), the KafkaBasedLog worker thread should catch TimeoutException and log a warning, which can occur if brokers are unavailable, otherwise the worker thread terminates. - Includes an enhancement to MockConsumer that allows simulating exceptions not just when polling but also when querying for offsets, which is necessary for testing the fix. --- .../kafka/clients/consumer/MockConsumer.java | 29 +++++-- .../kafka/connect/util/KafkaBasedLog.java | 4 + .../kafka/connect/util/KafkaBasedLogTest.java | 76 ++++++++++++++++++- .../internals/GlobalStreamThreadTest.java | 2 +- .../internals/StoreChangelogReaderTest.java | 2 +- .../processor/internals/StreamThreadTest.java | 2 +- 6 files changed, 103 insertions(+), 12 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 ce6a60bac2ab4..abeac096548ad 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 @@ -58,7 +58,8 @@ public class MockConsumer implements Consumer { private final Set paused; private Map>> records; - private KafkaException exception; + private KafkaException pollException; + private KafkaException offsetsException; private AtomicBoolean wakeup; private boolean closed; @@ -71,7 +72,7 @@ public MockConsumer(OffsetResetStrategy offsetResetStrategy) { this.beginningOffsets = new HashMap<>(); this.endOffsets = new HashMap<>(); this.pollTasks = new LinkedList<>(); - this.exception = null; + this.pollException = null; this.wakeup = new AtomicBoolean(false); this.committed = new HashMap<>(); } @@ -170,9 +171,9 @@ public synchronized ConsumerRecords poll(final Duration timeout) { throw new WakeupException(); } - if (exception != null) { - RuntimeException exception = this.exception; - this.exception = null; + if (pollException != null) { + RuntimeException exception = this.pollException; + this.pollException = null; throw exception; } @@ -213,8 +214,8 @@ public synchronized void addRecord(ConsumerRecord record) { recs.add(record); } - public synchronized void setException(KafkaException exception) { - this.exception = exception; + public synchronized void setPollException(KafkaException pollException) { + this.pollException = pollException; } @Override @@ -388,6 +389,11 @@ public synchronized Map offsetsForTimes(Map< @Override public synchronized Map beginningOffsets(Collection partitions) { + if (offsetsException != null) { + RuntimeException exception = this.offsetsException; + this.offsetsException = null; + throw exception; + } Map result = new HashMap<>(); for (TopicPartition tp : partitions) { Long beginningOffset = beginningOffsets.get(tp); @@ -400,6 +406,11 @@ public synchronized Map beginningOffsets(Collection endOffsets(Collection partitions) { + if (offsetsException != null) { + RuntimeException exception = this.offsetsException; + this.offsetsException = null; + throw exception; + } Map result = new HashMap<>(); for (TopicPartition tp : partitions) { Long endOffset = getEndOffset(endOffsets.get(tp)); @@ -410,6 +421,10 @@ public synchronized Map endOffsets(Collection endOffsets = new HashMap<>(); + endOffsets.put(TP0, 0L); + endOffsets.put(TP1, 0L); + consumer.updateEndOffsets(endOffsets); + store.start(); + final AtomicBoolean getInvoked = new AtomicBoolean(false); + final FutureCallback readEndFutureCallback = new FutureCallback<>(new Callback() { + @Override + public void onCompletion(Throwable error, Void result) { + getInvoked.set(true); + } + }); + consumer.schedulePollTask(new Runnable() { + @Override + public void run() { + // Once we're synchronized in a poll, start the read to end and schedule the exact set of poll events + // that should follow. This readToEnd call will immediately wakeup this consumer.poll() call without + // returning any data. + Map newEndOffsets = new HashMap<>(); + newEndOffsets.put(TP0, 1L); + newEndOffsets.put(TP1, 1L); + consumer.updateEndOffsets(newEndOffsets); + // Set exception to occur when getting offsets to read log to end. It'll be caught in the work thread, + // which will retry and eventually get the correct offsets and read log to end. + consumer.setOffsetsException(new TimeoutException("Failed to get offsets by times")); + store.readToEnd(readEndFutureCallback); + + // Should keep polling until it reaches current log end offset for all partitions + consumer.scheduleNopPollTask(); + consumer.scheduleNopPollTask(); + consumer.schedulePollTask(new Runnable() { + @Override + public void run() { + consumer.addRecord(new ConsumerRecord<>(TOPIC, 0, 0, 0L, TimestampType.CREATE_TIME, 0L, 0, 0, TP0_KEY, TP0_VALUE)); + consumer.addRecord(new ConsumerRecord<>(TOPIC, 1, 0, 0L, TimestampType.CREATE_TIME, 0L, 0, 0, TP0_KEY, TP0_VALUE_NEW)); + } + }); + + consumer.schedulePollTask(new Runnable() { + @Override + public void run() { + finishedLatch.countDown(); + } + }); + } + }); + readEndFutureCallback.get(10000, TimeUnit.MILLISECONDS); + assertTrue(getInvoked.get()); + assertTrue(finishedLatch.await(10000, TimeUnit.MILLISECONDS)); + assertEquals(CONSUMER_ASSIGNMENT, consumer.assignment()); + assertEquals(1L, consumer.position(TP0)); + + store.stop(); + + assertFalse(Whitebox.getInternalState(store, "thread").isAlive()); + assertTrue(consumer.closed()); + PowerMock.verifyAll(); + } + @Test public void testProducerError() throws Exception { expectStart(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java index c0e0de314964f..b67b664b621bf 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java @@ -236,7 +236,7 @@ public void shouldDieOnInvalidOffsetException() throws Exception { 10 * 1000, "Input record never consumed"); - mockConsumer.setException(new InvalidOffsetException("Try Again!") { + mockConsumer.setPollException(new InvalidOffsetException("Try Again!") { @Override public Set partitions() { return Collections.singleton(topicPartition); 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 b13ed535d625f..bf59c21f12c15 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 @@ -158,7 +158,7 @@ public void shouldRestoreAllMessagesFromBeginningWhenCheckpointNull() { public void shouldRecoverFromInvalidOffsetExceptionAndFinishRestore() { final int messages = 10; setupConsumer(messages, topicPartition); - consumer.setException(new InvalidOffsetException("Try Again!") { + consumer.setPollException(new InvalidOffsetException("Try Again!") { @Override public Set partitions() { return Collections.singleton(topicPartition); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index 11485e4e0c9d0..1a7c81e72b24d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -1234,7 +1234,7 @@ public void shouldRecoverFromInvalidOffsetExceptionOnRestoreAndFinishRestore() t () -> mockRestoreConsumer.position(changelogPartition) == 1L, "Never restore first record"); - mockRestoreConsumer.setException(new InvalidOffsetException("Try Again!") { + mockRestoreConsumer.setPollException(new InvalidOffsetException("Try Again!") { @Override public Set partitions() { return changelogPartitionSet; From 43c4fd9505202bd7ead61513212e16da3a211a77 Mon Sep 17 00:00:00 2001 From: pgwhalen Date: Sun, 17 Feb 2019 20:10:42 -0600 Subject: [PATCH 2/5] Fix GlobalStateManagerImplTest after MockConsumer method rename --- .../streams/processor/internals/GlobalStateManagerImplTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java index 809a0b679d331..2d29916ac1b57 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java @@ -290,7 +290,7 @@ public void shouldRestoreRecordsUpToHighwatermark() { @Test public void shouldRecoverFromInvalidOffsetExceptionAndRestoreRecords() { initializeConsumer(2, 0, t1); - consumer.setException(new InvalidOffsetException("Try Again!") { + consumer.setPollException(new InvalidOffsetException("Try Again!") { public Set partitions() { return Collections.singleton(t1); } From 1c8285ae751a0ffafe50fb71a3fee6eb92006a6b Mon Sep 17 00:00:00 2001 From: pgwhalen Date: Thu, 16 May 2019 20:39:59 -0500 Subject: [PATCH 3/5] Code review tweaks - Enhance logging on TimeoutException - Move setOffsetsException method to sensible place in MockConsumer --- .../org/apache/kafka/clients/consumer/MockConsumer.java | 8 ++++---- .../java/org/apache/kafka/connect/util/KafkaBasedLog.java | 3 ++- 2 files changed, 6 insertions(+), 5 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 abeac096548ad..8b8053823634f 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 @@ -218,6 +218,10 @@ public synchronized void setPollException(KafkaException pollException) { this.pollException = pollException; } + public synchronized void setOffsetsException(KafkaException exception) { + this.offsetsException = exception; + } + @Override public synchronized void commitAsync(Map offsets, OffsetCommitCallback callback) { ensureNotClosed(); @@ -421,10 +425,6 @@ public synchronized Map endOffsets(Collection Date: Thu, 23 May 2019 22:09:23 -0500 Subject: [PATCH 4/5] Preserve setException for compatibility --- .../apache/kafka/clients/consumer/MockConsumer.java | 12 ++++++++++-- 1 file changed, 10 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 8b8053823634f..ec356ca7fa98f 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 @@ -214,8 +214,16 @@ public synchronized void addRecord(ConsumerRecord record) { recs.add(record); } - public synchronized void setPollException(KafkaException pollException) { - this.pollException = pollException; + /** + * @deprecated Use {@link #setPollException(KafkaException)} instead + */ + @Deprecated + public synchronized void setException(KafkaException exception) { + this.pollException = exception; + } + + public synchronized void setPollException(KafkaException exception) { + this.pollException = exception; } public synchronized void setOffsetsException(KafkaException exception) { From 045131130aec48970581bcbb7353cdb01afffe74 Mon Sep 17 00:00:00 2001 From: pgwhalen Date: Thu, 23 May 2019 22:10:37 -0500 Subject: [PATCH 5/5] Call setException rather than replicate --- .../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 ec356ca7fa98f..d9367800b4222 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 @@ -219,7 +219,7 @@ public synchronized void addRecord(ConsumerRecord record) { */ @Deprecated public synchronized void setException(KafkaException exception) { - this.pollException = exception; + setPollException(exception); } public synchronized void setPollException(KafkaException exception) {