diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java index 76dae3e7f3d..ec08b4ad412 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java @@ -474,7 +474,7 @@ class WorkerSinkTask extends WorkerTask { private ConsumerRecords pollConsumer(long timeoutMs) { ConsumerRecords msgs = consumer.poll(Duration.ofMillis(timeoutMs)); - // Exceptions raised from the task during a rebalance should be rethrown to stop the worker + // Exceptions raised from the task during a rebalance should be rethrown to stop the task and mark it as failed if (rebalanceException != null) { RuntimeException e = rebalanceException; rebalanceException = null;