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..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 @@ -581,35 +581,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."); } }