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

mjsax 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 1ba3ca579c0 MINOR: Task-offset timer should not be reset if nothing 
was sent (#22715)
1ba3ca579c0 is described below

commit 1ba3ca579c0758b2f706721b6ea792c292b96c94
Author: Matthias J. Sax <[email protected]>
AuthorDate: Wed Jul 1 08:31:30 2026 -0700

    MINOR: Task-offset timer should not be reset if nothing was sent (#22715)
    
    Reviewers: Bill Bejeck <[email protected]>
---
 .../StreamsGroupHeartbeatRequestManager.java       | 13 +++---
 .../StreamsGroupHeartbeatRequestManagerTest.java   | 47 ++++++++++++++++++++++
 2 files changed, 54 insertions(+), 6 deletions(-)

diff --git 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index 361e313c104..bd554051f30 100644
--- 
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++ 
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -179,23 +179,24 @@ public class StreamsGroupHeartbeatRequestManager 
implements RequestManager {
                 final Map<StreamsRebalanceData.TaskId, Long> taskOffsetSum = 
streamsRebalanceData.taskOffsetSum();
                 final Map<StreamsRebalanceData.TaskId, Long> taskEndOffsetSum 
= streamsRebalanceData.taskEndOffsetSum();
 
+                final long now = time.milliseconds();
                 if (assignmentChanged
-                    || taskOffsetIntervalPassed()
+                    || taskOffsetIntervalPassed(now)
                     || 
hasAtLeastOneHotWarmupTask(reconciledAssignment.warmupTasks(), taskOffsetSum, 
taskEndOffsetSum)
                 ) {
                     // Task offsets and end-offsets are reported 
independently. A null field means "unchanged since the
                     // last heartbeat", so we send each one only when its 
value actually changed and leave it null
-                    // otherwise. reset() clears the snapshot on any 
error/disconnect, forcing a full resend afterwards.
+                    // otherwise. reset() clears the snapshot on any 
error/disconnect, forcing a full resend afterward.
                     if (!taskOffsetSum.equals(lastSentFields.taskOffsets)) {
                         data.setTaskOffsets(convertToList(taskOffsetSum));
                         lastSentFields.taskOffsets = taskOffsetSum;
+                        lastTaskOffsetIntervalTs = now;
                     }
                     if 
(!taskEndOffsetSum.equals(lastSentFields.taskEndOffsets)) {
                         
data.setTaskEndOffsets(convertToList(taskEndOffsetSum));
                         lastSentFields.taskEndOffsets = taskEndOffsetSum;
+                        lastTaskOffsetIntervalTs = now;
                     }
-
-                    lastTaskOffsetIntervalTs = time.milliseconds();
                 }
             }
             
data.setShutdownApplication(streamsRebalanceData.shutdownRequested());
@@ -211,8 +212,8 @@ public class StreamsGroupHeartbeatRequestManager implements 
RequestManager {
                 .collect(Collectors.toList());
         }
 
-        private boolean taskOffsetIntervalPassed() {
-            return lastTaskOffsetIntervalTs + 
streamsRebalanceData.taskOffsetIntervalMs() <= time.milliseconds();
+        private boolean taskOffsetIntervalPassed(final long now) {
+            return lastTaskOffsetIntervalTs + 
streamsRebalanceData.taskOffsetIntervalMs() <= now;
         }
 
         private boolean hasAtLeastOneHotWarmupTask(
diff --git 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
index cf6c661e646..7a9c5ff86c6 100644
--- 
a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
+++ 
b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManagerTest.java
@@ -1166,6 +1166,53 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertEquals(150L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
     }
 
+    @Test
+    public void testTaskOffsetIntervalNotAdvancedWhenNothingSent() {
+        // Entering the offset-send block via a non-interval trigger (here an 
assignment change)
+        // while the offsets are unchanged must NOT advance the 
task-offset-interval timer, because
+        // nothing was actually sent. If it did, the next interval-triggered 
resend of *changed*
+        // offsets would be withheld until a full interval after the spurious 
bump instead of a full
+        // interval after the last actual send.
+        final StreamsRebalanceData.TaskId task = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> offsets =
+            new AtomicReference<>(Map.of(task, 100L));
+        final StreamsRebalanceData rebalanceData = 
newRebalanceDataWithStandbyOffsets(
+            task,
+            offsets,
+            new AtomicReference<>(Map.of(task, 200L))
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // T=0: first build sends the offsets (assignment-changed trigger); 
the interval timer starts at 0.
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertEquals(100L, first.taskOffsets().get(0).offset());
+
+        // T=500 (mid-interval): a new assignment change re-enters the send 
block, but the offsets are
+        // unchanged, so nothing is sent. The interval timer must stay at 0 
(not advance to 500).
+        time.sleep(500);
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(task), // moved from standby to active so the assignment 
differs -> assignmentChanged
+            Set.of(),
+            Set.of(),     // no warmups, so hasAtLeastOneHotWarmupTask cannot 
be a trigger later
+            true
+        ));
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+
+        // T=1000: the offset changed and exactly one interval has elapsed 
since the last actual send
+        // (T=0). With the timer correctly still at 0, the interval trigger 
fires and the new offset is
+        // sent. With the bug (timer advanced to 500 at T=500), the interval 
would not be considered
+        // elapsed until T=1500 and the changed offset would be withheld.
+        time.sleep(500);
+        offsets.set(Map.of(task, 150L));
+        final StreamsGroupHeartbeatRequestData third = 
heartbeatState.buildRequestData();
+        assertNotNull(third.taskOffsets());
+        assertEquals(150L, third.taskOffsets().get(0).offset());
+    }
+
     @Test
     public void testTaskOffsetsAndEndOffsetsReportedIndependently() {
         // A null field means "unchanged", and the two fields are independent: 
one may be sent

Reply via email to