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 e14509da5f1e6..8a9aa268c5412 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 @@ -306,7 +306,7 @@ private synchronized boolean ensureCoordinatorReady(final Timer timer, boolean d if (future.isRetriable()) { log.debug("Coordinator discovery failed, refreshing metadata", future.exception()); timer.sleep(retryBackoff.backoff(attempts++)); - client.awaitMetadataUpdate(timer); + client.awaitMetadataUpdate(timer, disableWakeup); } else { fatalException = future.exception(); log.info("FindCoordinator request hit fatal exception", fatalException); diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerNetworkClient.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerNetworkClient.java index 1f07bbcaa2f9a..9096d49e90e96 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerNetworkClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerNetworkClient.java @@ -161,9 +161,20 @@ public boolean hasReadyNodes(long now) { * @return true if update succeeded, false otherwise. */ public boolean awaitMetadataUpdate(Timer timer) { + return awaitMetadataUpdate(timer, false); + } + + /** + * Block waiting on the metadata refresh with a timeout. + * + * @param timer Timer bounding how long this method can block + * @param disableWakeup true if we should not check for wakeups, false otherwise + * @return true if update succeeded, false otherwise. + */ + public boolean awaitMetadataUpdate(Timer timer, boolean disableWakeup) { int version = this.metadata.requestUpdate(false); do { - poll(timer); + poll(timer, null, disableWakeup); } while (this.metadata.updateVersion() == version && timer.notExpired()); return this.metadata.updateVersion() > version; } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java index 5c925d4821a28..343130da51217 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java @@ -323,6 +323,21 @@ public void testWakeupFromEnsureCoordinatorReady() { ); } + @Test + public void testNoWakeupFromAsyncCoordinatorReadyOnRetriableError() { + setupCoordinator(); + mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.COORDINATOR_NOT_AVAILABLE)); + consumerClient.wakeup(); + // The async variation should not throw WakeupException + coordinator.ensureCoordinatorReadyAsync(); + + consumerClient.wakeup(); + mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.COORDINATOR_NOT_AVAILABLE)); + assertThrows(WakeupException.class, () -> + coordinator.ensureCoordinatorReady(mockTime.timer(0)) + ); + } + @Test public void testTimeoutAndRetryJoinGroupIfNeeded() throws Exception { setupCoordinator();