diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java index 34f81db306d33..ab3c33beb876e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinator.java @@ -256,8 +256,8 @@ protected synchronized boolean ensureCoordinatorReady(final Timer timer) { log.debug("Coordinator discovery failed, refreshing metadata", future.exception()); client.awaitMetadataUpdate(timer); } else { - log.info("FindCoordinator request hit fatal exception", fatalException); fatalException = future.exception(); + log.info("FindCoordinator request hit fatal exception", fatalException); } } else if (coordinator != null && client.isUnavailable(coordinator)) { // we found the coordinator, but the connection has failed, so mark @@ -267,7 +267,7 @@ protected synchronized boolean ensureCoordinatorReady(final Timer timer) { } clearFindCoordinatorFuture(); - if (fatalException != null) + if (fatalException != null) throw fatalException; } while (coordinatorUnknown() && timer.notExpired());