This is an automated email from the ASF dual-hosted git repository.

AndrewJSchofield pushed a commit to branch 4.3
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/4.3 by this push:
     new 6ba09130a9c KAFKA-20253: Backport heartbeat CPU spin fixes (4.3) 
(#22964)
6ba09130a9c is described below

commit 6ba09130a9c509b41ea7a3ec5bc308b21cf1ac87
Author: Mingi Cho <[email protected]>
AuthorDate: Thu Jul 30 19:48:20 2026 +0900

    KAFKA-20253: Backport heartbeat CPU spin fixes (4.3) (#22964)
    
    Backports the following fixes to the `4.3` branch:
    
    - #22073
    - #22836
    
    #22073 applied cleanly.
    
    #22836 conflicted only in the import section of
    `ShareHeartbeatRequestManagerTest`. The 4.3 branch already uses
    `anyBoolean()`, so the existing `anyBoolean` import was retained.
    The remaining changes and regression tests were applied unchanged.
    
    Testing: `./gradlew :clients:test --tests
    'org.apache.kafka.clients.consumer.internals.AbstractCoordinatorTest'
    --tests
    
    
'org.apache.kafka.clients.consumer.internals.ConsumerHeartbeatRequestManagerTest'
    --tests
    
    'org.apache.kafka.clients.consumer.internals.CoordinatorRequestManagerTest'
    --tests
    
    
'org.apache.kafka.clients.consumer.internals.ShareHeartbeatRequestManagerTest'`
    
    All four test classes pass locally on this branch.
    
    Reviewers: Mickael Maison <[email protected]>
    
    ---------
    
    Co-authored-by: Evan Zhou <[email protected]>
---
 .../consumer/internals/AbstractCoordinator.java    |  1 +
 .../internals/AbstractHeartbeatRequestManager.java | 13 ++++-
 .../consumer/internals/CommitRequestManager.java   |  9 ++++
 .../internals/CoordinatorRequestManager.java       |  7 +++
 .../internals/AbstractCoordinatorTest.java         | 61 ++++++++++++++++++++++
 .../ConsumerHeartbeatRequestManagerTest.java       | 19 +++++++
 .../internals/CoordinatorRequestManagerTest.java   | 22 ++++++++
 .../ShareHeartbeatRequestManagerTest.java          | 18 +++++++
 8 files changed, 149 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 e14509da5f1..80d405e3bf0 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
@@ -1561,6 +1561,7 @@ public abstract class AbstractCoordinator implements 
Closeable {
             } catch (AuthenticationException e) {
                 log.error("An authentication error occurred in the heartbeat 
thread", e);
                 setFailureCause(e);
+                requestRejoin("authentication error in heartbeat thread");
             } catch (GroupAuthorizationException e) {
                 log.error("A group authorization error occurred in the 
heartbeat thread", e);
                 setFailureCause(e);
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 8392e5032f6..02be5779375 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 00c5812a5fa..e209188a93c 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 2c9c72e0520..3aaf7a2b8a2 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/AbstractCoordinatorTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/AbstractCoordinatorTest.java
index 5c925d4821a..a067cca8809 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
@@ -1288,6 +1288,67 @@ public class AbstractCoordinatorTest {
         }
     }
 
+    @Test
+    public void testAuthenticationErrorInHeartbeatThreadTriggersRejoin() 
throws Exception {
+        setupCoordinator();
+
+        mockClient.prepareResponse(groupCoordinatorResponse(node, 
Errors.NONE));
+        mockClient.prepareResponse(joinGroupFollowerResponse(1, memberId, 
leaderId, Errors.NONE));
+        mockClient.prepareResponse(syncGroupResponse(Errors.NONE));
+
+        final AuthenticationException authError = new 
AuthenticationException("test auth failure");
+        // A matcher exception leaves the response queued, so fail only the 
first heartbeat.
+        final AtomicBoolean authErrorThrown = new AtomicBoolean(false);
+
+        mockClient.prepareResponse(body -> {
+            if (!(body instanceof HeartbeatRequest))
+                return false;
+            if (authErrorThrown.compareAndSet(false, true))
+                throw authError;
+            return true;
+        }, heartbeatResponse(Errors.NONE));
+
+        coordinator.ensureActiveGroup();
+        assertFalse(coordinator.rejoinNeededOrPending(),
+            "Sanity: group should be stable before the auth error");
+
+        mockTime.sleep(HEARTBEAT_INTERVAL_MS);
+
+        TestUtils.waitForCondition(() -> {
+            try {
+                coordinator.pollHeartbeat(mockTime.milliseconds());
+                return false;
+            } catch (AuthenticationException e) {
+                assertSame(authError, e);
+                return true;
+            }
+        }, 3000, "HeartbeatThread did not propagate the authentication error 
in time");
+
+        assertTrue(coordinator.rejoinNeededOrPending(),
+            "Expected the heartbeat thread to request a rejoin after an 
AuthenticationException " +
+                "so the next poll() restarts the heartbeat machinery via 
ensureActiveGroup()");
+
+        mockClient.prepareResponse(joinGroupFollowerResponse(2, memberId, 
leaderId, Errors.NONE));
+        mockClient.prepareResponse(syncGroupResponse(Errors.NONE));
+
+        coordinator.ensureActiveGroup();
+
+        assertFalse(coordinator.rejoinNeededOrPending(), "Coordinator should 
have rejoined the group");
+        assertEquals(2, coordinator.generation().generationId);
+
+        mockClient.prepareResponse(body -> body instanceof HeartbeatRequest, 
heartbeatResponse(Errors.NONE));
+        mockTime.sleep(HEARTBEAT_INTERVAL_MS);
+
+        // Ensure the response was consumed instead of passing before a 
heartbeat was sent.
+        TestUtils.waitForCondition(() -> {
+            coordinator.pollHeartbeat(mockTime.milliseconds());
+            return !mockClient.hasPendingResponses() && 
!coordinator.heartbeat().hasInflight();
+        }, 3000, "Heartbeat was not sent and completed after recovering from 
the authentication error");
+
+        assertFalse(coordinator.rejoinNeededOrPending(),
+            "A successful post-recovery heartbeat should not trigger another 
rejoin");
+    }
+
     @Test
     public void testPollHeartbeatAwakesHeartbeatThread() throws Exception {
         final int longRetryBackoffMs = 10000;
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 906ae79518f..f7cfb31c262 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
@@ -341,6 +341,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 testTimerNotDue() {
         time.sleep(100); // time elapsed < heartbeatInterval, no heartbeat 
should be sent
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 c22af52cf89..af95f0c1cdd 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
@@ -236,6 +236,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 ce457a3a028..c41702a281e 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
@@ -167,6 +167,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