diff --git a/tools/src/main/java/org/apache/kafka/trogdor/common/WorkerUtils.java b/tools/src/main/java/org/apache/kafka/trogdor/common/WorkerUtils.java index 3d4871a17296b..978ca5ff5f4d5 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/common/WorkerUtils.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/common/WorkerUtils.java @@ -59,7 +59,7 @@ public final class WorkerUtils { * @param doneFuture The TaskWorker's doneFuture * @throws KafkaException A wrapped version of the exception. */ - public static void abort(Logger log, String what, Throwable exception, + public static void abortAndThrow(Logger log, String what, Throwable exception, KafkaFutureImpl doneFuture) throws KafkaException { log.warn("{} caught an exception: ", what, exception); doneFuture.complete(exception.getMessage()); diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressWorker.java index d85effc2833cb..6b35ed7a145ac 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConnectionStressWorker.java @@ -160,7 +160,7 @@ public void run() { } } } catch (Exception e) { - WorkerUtils.abort(log, "ConnectionStressRunnable", e, doneFuture); + WorkerUtils.abortAndThrow(log, "ConnectionStressRunnable", e, doneFuture); } } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java index 1e80209cc9fa2..68b302a8ed9a9 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java @@ -105,7 +105,7 @@ public void run() { } executor.submit(new CloseStatusUpdater(consumeTasks)); } catch (Throwable e) { - WorkerUtils.abort(log, "Prepare", e, doneFuture); + WorkerUtils.abortAndThrow(log, "Prepare", e, doneFuture); } } @@ -262,7 +262,8 @@ public Void call() throws Exception { startBatchMs = Time.SYSTEM.milliseconds(); } } catch (Exception e) { - WorkerUtils.abort(log, "ConsumeRecords", e, doneFuture); + consumer.close(); + WorkerUtils.abortAndThrow(log, "ConsumeRecords", e, doneFuture); } finally { statusUpdaterFuture.cancel(false); StatusData statusData = @@ -271,8 +272,8 @@ 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(); + doneFuture.complete(""); return null; } } @@ -311,7 +312,7 @@ public void run() { try { update(); } catch (Exception e) { - WorkerUtils.abort(log, "ConsumeStatusUpdater", e, doneFuture); + WorkerUtils.abortAndThrow(log, "ConsumeStatusUpdater", e, doneFuture); } } @@ -343,7 +344,7 @@ public void run() { try { update(); } catch (Exception e) { - WorkerUtils.abort(log, "ConsumeStatusUpdater", e, doneFuture); + WorkerUtils.abortAndThrow(log, "ConsumeStatusUpdater", e, doneFuture); } } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchWorker.java index 84b94d582477e..12b4b9dfc4b05 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ProduceBenchWorker.java @@ -123,7 +123,7 @@ public void run() { status.update(new TextNode("Created " + newTopics.keySet().size() + " topic(s)")); executor.submit(new SendRecords(active)); } catch (Throwable e) { - WorkerUtils.abort(log, "Prepare", e, doneFuture); + WorkerUtils.abortAndThrow(log, "Prepare", e, doneFuture); } } } @@ -248,7 +248,7 @@ public Void call() throws Exception { producer.close(); } } catch (Exception e) { - WorkerUtils.abort(log, "SendRecords", e, doneFuture); + WorkerUtils.abortAndThrow(log, "SendRecords", e, doneFuture); } finally { statusUpdaterFuture.cancel(false); StatusData statusData = new StatusUpdater(histogram, transactionsCommitted).update(); @@ -315,7 +315,7 @@ public void run() { try { update(); } catch (Exception e) { - WorkerUtils.abort(log, "StatusUpdater", e, doneFuture); + WorkerUtils.abortAndThrow(log, "StatusUpdater", e, doneFuture); } } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java index 669fafcc75ed5..3fa83ae47402c 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java @@ -153,7 +153,7 @@ public void run() { executor.scheduleWithFixedDelay( new StatusUpdater(), 30, 30, TimeUnit.SECONDS); } catch (Throwable e) { - WorkerUtils.abort(log, "Prepare", e, doneFuture); + WorkerUtils.abortAndThrow(log, "Prepare", e, doneFuture); } } } @@ -262,7 +262,7 @@ public void onCompletion(RecordMetadata metadata, Exception exception) { }); } } catch (Throwable e) { - WorkerUtils.abort(log, "ProducerRunnable", e, doneFuture); + WorkerUtils.abortAndThrow(log, "ProducerRunnable", e, doneFuture); } finally { log.info("{}: ProducerRunnable is exiting. messagesSent={}; uniqueMessagesSent={}; " + "ackedSends={}.", id, messagesSent, uniqueMessagesSent, @@ -368,7 +368,7 @@ public void run() { } } } catch (Throwable e) { - WorkerUtils.abort(log, "ConsumerRunnable", e, doneFuture); + WorkerUtils.abortAndThrow(log, "ConsumerRunnable", e, doneFuture); } finally { log.info("{}: ConsumerRunnable is exiting. Invoked poll {} time(s). " + "messagesReceived = {}; uniqueMessagesReceived = {}.", @@ -383,7 +383,7 @@ public void run() { try { update(); } catch (Exception e) { - WorkerUtils.abort(log, "StatusUpdater", e, doneFuture); + WorkerUtils.abortAndThrow(log, "StatusUpdater", e, doneFuture); } }