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());

Reply via email to