From 2d0848b0ebb052a94f719b06aa4ee2547e5ce717 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 8 Sep 2022 17:12:03 -0700 Subject: [PATCH 1/3] should not wake-up with non-blocking coordinator discovery --- .../internals/AbstractCoordinator.java | 9 ++++- .../internals/AbstractCoordinatorTest.java | 33 +++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) 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 d2ece9efc587c..11894ad9a3713 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 @@ -249,7 +249,14 @@ protected synchronized boolean ensureCoordinatorReady(final Timer timer) { throw fatalException; } final RequestFuture future = lookupCoordinator(); - client.poll(future, timer); + + // if we do not want to block on discovering coordinator at all, + // then we should not try to poll in a loop, and should not throw wake-up exception either + if (timer.timeoutMs() == 0L) { + client.poll(timer, future, true); + } else { + client.poll(future, timer); + } if (!future.isDone()) { // ran out of time 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 cbc4e7495e161..28898c3ebff30 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 @@ -272,6 +272,39 @@ public void testCoordinatorDiscoveryBackoff() { assertTrue(endTime - initialTime >= RETRY_BACKOFF_MS); } + @Test + public void testNoWakeupWhenNonBlockingDiscoverCoordinator() { + setupCoordinator(); + + mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.NONE)); + + consumerClient.wakeup(); + + coordinator.ensureCoordinatorReady(mockTime.timer(0)); + + // a follow-up poll should still throw + try { + coordinator.joinGroupIfNeeded(mockTime.timer(0)); + fail("Should have woken up from joinGroupIfNeeded()"); + } catch (WakeupException ignored) { + } + } + + @Test + public void testWakeupWhenBlockingDiscoverCoordinator() throws Exception { + setupCoordinator(); + + mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.NONE)); + + consumerClient.wakeup(); + + try { + coordinator.ensureCoordinatorReady(mockTime.timer(1)); + fail("Should have woken up from ensureCoordinatorReady()"); + } catch (WakeupException ignored) { + } + } + @Test public void testTimeoutAndRetryJoinGroupIfNeeded() throws Exception { setupCoordinator(); From 56aa710e844ac0ebd8f35604e1e80ea9228e6f72 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 12 Sep 2022 13:02:31 -0700 Subject: [PATCH 2/3] Add explicit flag to disable wakeups in `ensureCoordinatorReady` --- .../internals/AbstractCoordinator.java | 26 +++++++++------ .../internals/ConsumerCoordinator.java | 13 +++++--- .../internals/ConsumerNetworkClient.java | 17 +++++++++- .../internals/AbstractCoordinatorTest.java | 32 ++++--------------- 4 files changed, 47 insertions(+), 41 deletions(-) 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 11894ad9a3713..c78caaba0596d 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 @@ -231,14 +231,27 @@ protected abstract void onJoinComplete(int generation, protected void onLeavePrepare() {} /** - * Visible for testing. - * * Ensure that the coordinator is ready to receive requests. * * @param timer Timer bounding how long this method can block * @return true If coordinator discovery and initial connection succeeded, false otherwise */ protected synchronized boolean ensureCoordinatorReady(final Timer timer) { + return ensureCoordinatorReady(timer, false); + } + + /** + * Ensure that the coordinator is ready to receive requests. This will return + * immediately without blocking. It is intended to be called in an asynchronous + * context when wakeups are not expected. + * + * @return true If coordinator discovery and initial connection succeeded, false otherwise + */ + protected synchronized boolean ensureCoordinatorReadyAsync() { + return ensureCoordinatorReady(time.timer(0), true); + } + + private synchronized boolean ensureCoordinatorReady(final Timer timer, boolean disableWakeup) { if (!coordinatorUnknown()) return true; @@ -249,14 +262,7 @@ protected synchronized boolean ensureCoordinatorReady(final Timer timer) { throw fatalException; } final RequestFuture future = lookupCoordinator(); - - // if we do not want to block on discovering coordinator at all, - // then we should not try to poll in a loop, and should not throw wake-up exception either - if (timer.timeoutMs() == 0L) { - client.poll(timer, future, true); - } else { - client.poll(future, timer); - } + client.poll(future, timer, disableWakeup); if (!future.isDone()) { // ran out of time diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java index 5228c60e0fb9a..f955a5fa59cd5 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ConsumerCoordinator.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.clients.consumer.internals; -import java.time.Duration; import java.util.SortedSet; import java.util.TreeSet; import org.apache.kafka.clients.GroupRebalanceConfig; @@ -489,10 +488,14 @@ void maybeUpdateSubscriptionMetadata() { } } - private boolean coordinatorUnknownAndUnready(Timer timer) { + private boolean coordinatorUnknownAndUnreadySync(Timer timer) { return coordinatorUnknown() && !ensureCoordinatorReady(timer); } + private boolean coordinatorUnknownAndUnreadyAsync() { + return coordinatorUnknown() && !ensureCoordinatorReadyAsync(); + } + /** * Poll for coordinator events. This ensures that the coordinator is known and that the consumer * has joined the group (if it is using group management). This also handles periodic offset commits @@ -518,7 +521,7 @@ public boolean poll(Timer timer, boolean waitForJoinGroup) { // Always update the heartbeat last poll time so that the heartbeat thread does not leave the // group proactively due to application inactivity even if (say) the coordinator cannot be found. pollHeartbeat(timer.currentTimeMs()); - if (coordinatorUnknownAndUnready(timer)) { + if (coordinatorUnknownAndUnreadySync(timer)) { return false; } @@ -1052,7 +1055,7 @@ public RequestFuture commitOffsetsAsync(final Map offsets, return true; do { - if (coordinatorUnknownAndUnready(timer)) { + if (coordinatorUnknownAndUnreadySync(timer)) { return false; } 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 4b9112016e986..6ba5666316554 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 @@ -211,8 +211,23 @@ public void poll(RequestFuture future) { * @throws InterruptException if the calling thread is interrupted */ public boolean poll(RequestFuture future, Timer timer) { + return poll(future, timer, false); + } + + /** + * Block until the provided request future request has finished or the timeout has expired. + * + * @param future The request future to wait for + * @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 the future is done, false otherwise + * @throws WakeupException if {@link #wakeup()} is called from another thread + * @throws InterruptException if the calling thread is interrupted + */ + public boolean poll(RequestFuture future, Timer timer, boolean disableWakeup) { do { - poll(timer, future); + poll(timer, future, disableWakeup); } while (!future.isDone() && timer.notExpired()); return future.isDone(); } 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 28898c3ebff30..4471b6f88cf2e 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 @@ -273,36 +273,18 @@ public void testCoordinatorDiscoveryBackoff() { } @Test - public void testNoWakeupWhenNonBlockingDiscoverCoordinator() { + public void testWakeupFromEnsureCoordinatorReady() { setupCoordinator(); - mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.NONE)); - consumerClient.wakeup(); - coordinator.ensureCoordinatorReady(mockTime.timer(0)); - - // a follow-up poll should still throw - try { - coordinator.joinGroupIfNeeded(mockTime.timer(0)); - fail("Should have woken up from joinGroupIfNeeded()"); - } catch (WakeupException ignored) { - } - } - - @Test - public void testWakeupWhenBlockingDiscoverCoordinator() throws Exception { - setupCoordinator(); + // No wakeup should occur from the async variation. + coordinator.ensureCoordinatorReadyAsync(); - mockClient.prepareResponse(groupCoordinatorResponse(node, Errors.NONE)); - - consumerClient.wakeup(); - - try { - coordinator.ensureCoordinatorReady(mockTime.timer(1)); - fail("Should have woken up from ensureCoordinatorReady()"); - } catch (WakeupException ignored) { - } + // But should wakeup in sync variation even if timer is 0. + assertThrows(WakeupException.class, () -> { + coordinator.ensureCoordinatorReady(mockTime.timer(0)); + }); } @Test From b1e0a5365f6dbfea89a7a3216884d1387ba6d491 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 12 Sep 2022 19:29:22 -0700 Subject: [PATCH 3/3] Document WakeupException throw conditions --- .../kafka/clients/consumer/internals/ConsumerNetworkClient.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 6ba5666316554..6646dc6c893e0 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 @@ -222,7 +222,7 @@ public boolean poll(RequestFuture future, Timer timer) { * @param disableWakeup true if we should not check for wakeups, false otherwise * * @return true if the future is done, false otherwise - * @throws WakeupException if {@link #wakeup()} is called from another thread + * @throws WakeupException if {@link #wakeup()} is called from another thread and `disableWakeup` is false * @throws InterruptException if the calling thread is interrupted */ public boolean poll(RequestFuture future, Timer timer, boolean disableWakeup) {