This is an automated email from the ASF dual-hosted git repository.
AndrewJSchofield pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 28de22de34e KAFKA-20253 Heartbeat CPU spin fix for Consumer protocol
(#22836)
28de22de34e is described below
commit 28de22de34eb5e93a803c37b46e3d0fa2fdf8980
Author: Evan Zhou <[email protected]>
AuthorDate: Thu Jul 16 04:26:25 2026 -0500
KAFKA-20253 Heartbeat CPU spin fix for Consumer protocol (#22836)
This PR fixes the heartbeat spin issue for the `consumer` protocol. In
the process of reproducing the issue, I found that both the application
and background threads have high CPU, and fixes for both threads are in
this PR
Jira: https://issues.apache.org/jira/browse/KAFKA-20253
Reviewers: Andrew Schofield <[email protected]>
---
.../internals/AbstractHeartbeatRequestManager.java | 13 ++++++++++++-
.../consumer/internals/CommitRequestManager.java | 9 +++++++++
.../internals/CoordinatorRequestManager.java | 7 +++++++
.../ConsumerHeartbeatRequestManagerTest.java | 19 +++++++++++++++++++
.../internals/CoordinatorRequestManagerTest.java | 22 ++++++++++++++++++++++
.../ShareHeartbeatRequestManagerTest.java | 19 +++++++++++++++++++
6 files changed, 88 insertions(+), 1 deletion(-)
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
index f98a9cbdad3..8226d39535a 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AbstractHeartbeatRequestManager.java
@@ -256,7 +256,18 @@ public abstract class AbstractHeartbeatRequestManager<R
extends AbstractResponse
if (membershipManager().state() == MemberState.UNSUBSCRIBED) {
return Long.MAX_VALUE;
}
- if (pollTimer.isExpired() || (membershipManager().shouldHeartbeatNow()
&& !heartbeatRequestState.requestInFlight())) {
+ if (pollTimer.isExpired()) {
+ return 0L;
+ }
+ // KAFKA-20253: mirror the guard in poll(). A heartbeat is only sent
when the coordinator is known
+ // and the member is in a state that heartbeats. When the coordinator
is unavailable (e.g. after a
+ // re-authentication failure) or the member should skip heartbeats
(FATAL/FENCED/STALE/UNSUBSCRIBED),
+ // poll() returns EMPTY, so falling through to the timer-based
branches below would return 0 (the
+ // heartbeat timer is left permanently expired) and busy-spin both the
application and network threads.
+ if (coordinatorRequestManager.coordinator().isEmpty() ||
membershipManager().shouldSkipHeartbeat()) {
+ return heartbeatRequestState.heartbeatIntervalMs();
+ }
+ if (membershipManager().shouldHeartbeatNow() &&
!heartbeatRequestState.requestInFlight()) {
return 0L;
}
return Math.min(pollTimer.remainingMs() / 2,
heartbeatRequestState.timeToNextHeartbeatMs(currentTimeMs));
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java
index 7912c0c4135..27fcd53607f 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CommitRequestManager.java
@@ -1532,6 +1532,15 @@ public class CommitRequestManager implements
RequestManager, MemberStateListener
public long remainingMs(final long currentTimeMs) {
this.timer.update(currentTimeMs);
+ // KAFKA-20253: If the auto-commit interval has elapsed but a
previous auto-commit is still
+ // in-flight (for example it cannot complete because the
coordinator is unavailable after a
+ // failed re-authentication), a new auto-commit cannot be started
yet. Returning 0 here would
+ // busy-spin the application thread, since this value feeds
AsyncKafkaConsumer.pollForFetches()
+ // via maximumTimeToWait(). Wait for the interval instead; the
network thread still wakes on the
+ // in-flight commit's response, which resets this timer.
+ if (this.timer.isExpired() && this.hasInflightCommit) {
+ return autoCommitInterval;
+ }
return this.timer.remainingMs();
}
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManager.java
index 2056ad66c23..0ec3504a725 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManager.java
@@ -107,6 +107,13 @@ public class CoordinatorRequestManager implements
RequestManager {
return new NetworkClientDelegate.PollResult(request);
}
+ // When a request is in flight, remainingBackoffMs() can be 0, and
returning 0 tells the network thread to
+ // poll again immediately which causes a busy spin. Wait instead by
returning a PollResult with a Long.MAX_VALUE
+ // backoff
+ if (coordinatorRequestState.requestInFlight()) {
+ return EMPTY;
+ }
+
return new
NetworkClientDelegate.PollResult(coordinatorRequestState.remainingBackoffMs(currentTimeMs));
}
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
index d9405a51929..b873571e00c 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ConsumerHeartbeatRequestManagerTest.java
@@ -306,6 +306,25 @@ public class ConsumerHeartbeatRequestManagerTest
}
}
+ /**
+ * KAFKA-20253: when the coordinator is unavailable (e.g. after a
re-authentication failure),
+ * poll() returns EMPTY, so no heartbeat can be sent. maximumTimeToWait()
must return a positive
+ * value in that case; returning 0 busy-spins the application thread (and,
via wakeups, the
+ * consumer network thread), which is the AsyncKafkaConsumer high-CPU loop
in this ticket.
+ */
+ @Test
+ public void testMaximumTimeToWaitWhenCoordinatorUnavailableDoesNotSpin() {
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.empty());
+ when(membershipManager.state()).thenReturn(MemberState.STABLE);
+ when(membershipManager.shouldHeartbeatNow()).thenReturn(true);
+
+ long result =
heartbeatRequestManager.maximumTimeToWait(time.milliseconds());
+
+ assertTrue(result > 0,
+ "maximumTimeToWait must be > 0 when the coordinator is unavailable
to avoid a busy-spin; got " + result);
+ assertEquals(DEFAULT_HEARTBEAT_INTERVAL_MS, result);
+ }
+
@Test
public void testHeartbeatNotSentIfAnotherOneInFlight() {
time.sleep(DEFAULT_HEARTBEAT_INTERVAL_MS);
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManagerTest.java
index 955e60d1a2d..7a76955d320 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/CoordinatorRequestManagerTest.java
@@ -238,6 +238,28 @@ public class CoordinatorRequestManagerTest {
assertEquals(1, res2.unsentRequests.size());
}
+ @Test
+ public void testNoBusyPollWhileFindCoordinatorRequestInFlight() {
+ // KAFKA-20253: while a FindCoordinator request is in flight and its
backoff has already
+ // elapsed, poll() must not return timeUntilNextPollMs == 0. Doing so
drives the consumer
+ // network thread into a NetworkClient.poll(0) busy-spin, since there
is nothing to send
+ // until the in-flight request completes.
+ CoordinatorRequestManager coordinatorManager =
setupCoordinatorManager(GROUP_ID);
+
+ // First poll sends a FindCoordinator request, marking it in-flight.
Do NOT complete it.
+ NetworkClientDelegate.PollResult res =
coordinatorManager.poll(time.milliseconds());
+ assertEquals(1, res.unsentRequests.size());
+
+ // Advance well past the retry backoff while the request is still in
flight.
+ time.sleep(60_000);
+
+ NetworkClientDelegate.PollResult res2 =
coordinatorManager.poll(time.milliseconds());
+ assertEquals(0, res2.unsentRequests.size(), "no new request should be
sent while one is in flight");
+ assertTrue(res2.timeUntilNextPollMs > 0,
+ "must not busy-poll (timeUntilNextPollMs == 0) while a
FindCoordinator request is in flight; got "
+ + res2.timeUntilNextPollMs);
+ }
+
@ParameterizedTest
@EnumSource(value = Errors.class, names = {"NONE",
"COORDINATOR_NOT_AVAILABLE"})
public void testClearFatalErrorWhenReceivingSuccessfulResponse(Errors
error) {
diff --git
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
index d1b02340619..0f125f09dfb 100644
---
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
+++
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareHeartbeatRequestManagerTest.java
@@ -62,6 +62,7 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -156,6 +157,24 @@ public class ShareHeartbeatRequestManagerTest
backgroundEventHandler);
}
+ /**
+ * ShareConsumerImpl runs on the same ConsumerNetworkThread as the
AsyncKafkaConsumer, thus we also test to ensure
+ * we don't enter a CPU spin state. Refer to KAFKA-20253
+ */
+ @Test
+ public void testMaximumTimeToWaitWhenHeartbeatShouldBeSkippedDoesNotSpin()
{
+
when(coordinatorRequestManager.coordinator()).thenReturn(Optional.of(new
Node(1, "localhost", 9999)));
+ when(membershipManager.state()).thenReturn(MemberState.FATAL);
+ when(membershipManager.shouldSkipHeartbeat()).thenReturn(true);
+
when(heartbeatRequestState.timeToNextHeartbeatMs(anyLong())).thenReturn(0L);
+
+ long result =
heartbeatRequestManager.maximumTimeToWait(time.milliseconds());
+
+ assertTrue(result > 0,
+ "maximumTimeToWait must be > 0 while heartbeats are skipped to
avoid a busy-spin; got " + result);
+ assertEquals(DEFAULT_HEARTBEAT_INTERVAL_MS, result);
+ }
+
@Test
public void testHeartbeatOnStartup() {
NetworkClientDelegate.PollResult result =
heartbeatRequestManager.poll(time.milliseconds());