From 591ae58d27a3243c0a3c9d6e197c58c118e68d35 Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 2 Jun 2020 14:11:42 -0500 Subject: [PATCH 1/2] KAFKA-10084: Fix EosTestDriver end offset --- .../kafka/streams/tests/EosTestDriver.java | 45 +++++++++---------- 1 file changed, 22 insertions(+), 23 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java b/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java index e95c3541e63ab..972c0b3b40109 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java @@ -44,6 +44,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.Iterator; import java.util.LinkedList; import java.util.List; @@ -51,6 +52,7 @@ import java.util.Map; import java.util.Properties; import java.util.Random; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -581,35 +583,32 @@ private static void verifyAllTransactionFinished(final KafkaConsumer topicEndOffsets; - - try (final KafkaConsumer consumerUncommitted = new KafkaConsumer<>(consumerProps)) { - topicEndOffsets = consumerUncommitted.endOffsets(partitions); - } final long maxWaitTime = System.currentTimeMillis() + MAX_IDLE_TIME_MS; - while (!topicEndOffsets.isEmpty() && System.currentTimeMillis() < maxWaitTime) { - consumer.seekToEnd(partitions); - - final Iterator iterator = partitions.iterator(); - while (iterator.hasNext()) { - final TopicPartition topicPartition = iterator.next(); - final long position = consumer.position(topicPartition); - - if (position == topicEndOffsets.get(topicPartition)) { - iterator.remove(); - topicEndOffsets.remove(topicPartition); - System.out.println("Removing " + topicPartition + " at position " + position); - } else if (consumer.position(topicPartition) > topicEndOffsets.get(topicPartition)) { - throw new IllegalStateException("Offset for partition " + topicPartition + " is larger than topic endOffset: " + position + " > " + topicEndOffsets.get(topicPartition)); - } else { - System.out.println("Retry " + topicPartition + " at position " + position); + try (final KafkaConsumer consumerUncommitted = new KafkaConsumer<>(consumerProps)) { + while (System.currentTimeMillis() < maxWaitTime) { + consumer.seekToEnd(partitions); + final Map topicEndOffsets = consumerUncommitted.endOffsets(partitions); + + final Iterator iterator = partitions.iterator(); + while (iterator.hasNext()) { + final TopicPartition topicPartition = iterator.next(); + final long position = consumer.position(topicPartition); + + if (position == topicEndOffsets.get(topicPartition)) { + iterator.remove(); + System.out.println("Removing " + topicPartition + " at position " + position); + } else if (consumer.position(topicPartition) > topicEndOffsets.get(topicPartition)) { + throw new IllegalStateException("Offset for partition " + topicPartition + " is larger than topic endOffset: " + position + " > " + topicEndOffsets.get(topicPartition)); + } else { + System.out.println("Retry " + topicPartition + " at position " + position); + } } + sleep(1000L); } - sleep(1000L); } - if (!topicEndOffsets.isEmpty()) { + if (!partitions.isEmpty()) { throw new RuntimeException("Could not read all verification records. Did not receive any new record within the last " + (MAX_IDLE_TIME_MS / 1000L) + " sec."); } } From c3b131076fffa8a9ee4d3aeb5dd53d1eb21f31bc Mon Sep 17 00:00:00 2001 From: John Roesler Date: Tue, 2 Jun 2020 15:02:31 -0500 Subject: [PATCH 2/2] fix style --- .../test/java/org/apache/kafka/streams/tests/EosTestDriver.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java b/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java index 972c0b3b40109..45843aac11a25 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/EosTestDriver.java @@ -44,7 +44,6 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; -import java.util.HashSet; import java.util.Iterator; import java.util.LinkedList; import java.util.List; @@ -52,7 +51,6 @@ import java.util.Map; import java.util.Properties; import java.util.Random; -import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit;