From 15d3e2b96ffbc7eedac61e5ac86b96cc3e3191ff Mon Sep 17 00:00:00 2001 From: Kowshik Prakasam Date: Mon, 31 May 2021 00:55:30 -0700 Subject: [PATCH] KAFKA-12867: Fix ConsumeBenchWorker exit behavior for maxMessages config --- .../org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/trogdor/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java b/trogdor/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java index f6067b59af4e9..84ce1d333f18b 100644 --- a/trogdor/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java +++ b/trogdor/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java @@ -282,7 +282,6 @@ public Void call() throws Exception { log.info("{} Consumed total number of messages={}, bytes={} in {} ms. status: {}", clientId, messagesConsumed, bytesConsumed, curTimeMs - startTimeMs, statusData); } - doneFuture.complete(""); consumer.close(); return null; } @@ -307,6 +306,7 @@ public void run() { } statusUpdaterFuture.cancel(false); statusUpdater.update(); + doneFuture.complete(""); } }