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 e78cd1fb7c7 KAFKA-20116: Send task-(end)-offset only if changed (4/N) 
(#22645)
e78cd1fb7c7 is described below

commit e78cd1fb7c77af7b6e63e8f1bea9120fb193508f
Author: Matthias J. Sax <[email protected]>
AuthorDate: Thu Jun 25 21:08:50 2026 -0700

    KAFKA-20116: Send task-(end)-offset only if changed (4/N) (#22645)
    
    KIP-1071 specifies, that any fields in the StreamsHeartbeatRequest
    should only be set if they did change since the last heartbeat. This PR
    add this behavior for task-offset-sum and task-end-offset-sum fields.
    
    Reviewers: Lucas Brutschy <[email protected]>
---
 checkstyle/suppressions.xml                        |   2 +-
 .../StreamsGroupHeartbeatRequestManager.java       |  29 ++-
 .../StreamsGroupHeartbeatRequestManagerTest.java   | 230 +++++++++++++++++++--
 3 files changed, 235 insertions(+), 26 deletions(-)

diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index cd064039853..4f252572fb2 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -54,7 +54,7 @@
     <suppress id="dontUseSystemExit"
               files="Exit.java"/>
     <suppress checks="ClassFanOutComplexity"
-              
files="(AbstractFetch|Sender|SenderTest|ConsumerCoordinator|KafkaConsumer|KafkaProducer|Utils|TransactionManager|TransactionManagerTest|KafkaAdminClient|NetworkClient|Admin|RaftClientTestContext|TestingMetricsInterceptingAdminClient|SaslServerAuthenticator|SaslAuthenticatorTest|Errors|AbstractRequest|AbstractResponse).java"/>
+              
files="(AbstractFetch|Sender|SenderTest|ConsumerCoordinator|KafkaConsumer|KafkaProducer|Utils|TransactionManager|TransactionManagerTest|KafkaAdminClient|NetworkClient|Admin|RaftClientTestContext|TestingMetricsInterceptingAdminClient|SaslServerAuthenticator|SaslAuthenticatorTest|Errors|AbstractRequest|AbstractResponse|StreamsGroupHeartbeatRequestManagerTest).java"/>
     <suppress checks="NPathComplexity"
               files="(SaslServerAuthenticator).java"/>
 
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 d59a2f2c90f..361e313c104 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
@@ -78,12 +78,16 @@ public class StreamsGroupHeartbeatRequestManager implements 
RequestManager {
         static class LastSentFields {
 
             private StreamsRebalanceData.Assignment assignment = null;
+            private Map<StreamsRebalanceData.TaskId, Long> taskOffsets = null;
+            private Map<StreamsRebalanceData.TaskId, Long> taskEndOffsets = 
null;
 
             LastSentFields() {
             }
 
             void reset() {
                 assignment = null;
+                taskOffsets = null;
+                taskEndOffsets = null;
             }
         }
 
@@ -151,8 +155,14 @@ public class StreamsGroupHeartbeatRequestManager 
implements RequestManager {
                 data.setActiveTasks(fromStreamsToHeartbeatRequest(Set.of()));
                 data.setStandbyTasks(fromStreamsToHeartbeatRequest(Set.of()));
                 data.setWarmupTasks(fromStreamsToHeartbeatRequest(Set.of()));
-                
data.setTaskOffsets(convertToList(streamsRebalanceData.taskOffsetSum()));
-                
data.setTaskEndOffsets(convertToList(streamsRebalanceData.taskEndOffsetSum()));
+                // call both methods only once, as they invoke an expensive 
`supplier`
+                final Map<StreamsRebalanceData.TaskId, Long> taskOffsetSum = 
streamsRebalanceData.taskOffsetSum();
+                final Map<StreamsRebalanceData.TaskId, Long> taskEndOffsetSum 
= streamsRebalanceData.taskEndOffsetSum();
+                data.setTaskOffsets(convertToList(taskOffsetSum));
+                data.setTaskEndOffsets(convertToList(taskEndOffsetSum));
+                // Record what we sent so the first non-joining heartbeat does 
not redundantly resend unchanged offsets.
+                lastSentFields.taskOffsets = taskOffsetSum;
+                lastSentFields.taskEndOffsets = taskEndOffsetSum;
                 lastTaskOffsetIntervalTs = time.milliseconds();
             } else {
                 final StreamsRebalanceData.Assignment reconciledAssignment = 
streamsRebalanceData.reconciledAssignment();
@@ -173,10 +183,17 @@ public class StreamsGroupHeartbeatRequestManager 
implements RequestManager {
                     || taskOffsetIntervalPassed()
                     || 
hasAtLeastOneHotWarmupTask(reconciledAssignment.warmupTasks(), taskOffsetSum, 
taskEndOffsetSum)
                 ) {
-
-                    // TODO: send only if changed this last time
-                    data.setTaskOffsets(convertToList(taskOffsetSum));
-                    data.setTaskEndOffsets(convertToList(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.
+                    if (!taskOffsetSum.equals(lastSentFields.taskOffsets)) {
+                        data.setTaskOffsets(convertToList(taskOffsetSum));
+                        lastSentFields.taskOffsets = taskOffsetSum;
+                    }
+                    if 
(!taskEndOffsetSum.equals(lastSentFields.taskEndOffsets)) {
+                        
data.setTaskEndOffsets(convertToList(taskEndOffsetSum));
+                        lastSentFields.taskEndOffsets = taskEndOffsetSum;
+                    }
 
                     lastTaskOffsetIntervalTs = time.milliseconds();
                 }
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 315c58c0b89..cf6c661e646 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
@@ -66,6 +66,7 @@ import java.util.Optional;
 import java.util.Properties;
 import java.util.Set;
 import java.util.UUID;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
@@ -856,14 +857,28 @@ class StreamsGroupHeartbeatRequestManagerTest {
     public void testHotWarmupTaskTriggersSendWhenLagAtOrBelowThreshold() {
         // A v1+ broker provides a positive acceptableRecoveryLag. When a 
warmup's lag
         // (endOffset - offset) is at or below the threshold, hasHotWarmupTask 
triggers
-        // an early send so the broker can promote the warmup promptly.
+        // an early send (before the task-offset interval elapses) so the 
broker can promote
+        // the warmup promptly. The send still happens only when the offset 
actually changed.
         final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
-        final StreamsRebalanceData rebalanceData = newRebalanceDataWithWarmup(
-            warmupTaskId,
-            900L, // offset
-            1000L, // endOffset → lag = 100
-            100L  // acceptableRecoveryLag
+        final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> offsets =
+            new AtomicReference<>(Map.of(warmupTaskId, 900L)); // lag = 1000 - 
900 = 100 → hot
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            offsets::get,
+            () -> Map.of(warmupTaskId, 1000L)
         );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(),
+            Set.of(warmupTaskId),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(100L);
 
         final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
@@ -878,13 +893,17 @@ class StreamsGroupHeartbeatRequestManagerTest {
         final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
         assertEquals(900L, first.taskOffsets().get(0).offset());
 
-        // Second STABLE build:
-        //  - assignmentChanged is false
-        //  - task.offset.interval.ms did not pass, as we did not advance time
-        // the only remaining candidate trigger is hasHotWarmupTask — and with 
valid acceptableRecoveryLag and low lag
-        // we expect the offset to be sent
+        // Second STABLE build without advancing time and with an unchanged 
offset:
+        //  - assignmentChanged is false, the interval did not pass, the 
warmup is still hot
+        // but since the offset did not change since the last heartbeat, 
nothing is resent.
         final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
-        assertEquals(900L, second.taskOffsets().get(0).offset());
+        assertNull(second.taskOffsets());
+
+        // The warmup makes progress (still hot). The hot-warmup trigger lets 
us report the new
+        // offset promptly, before the task-offset interval elapses.
+        offsets.set(Map.of(warmupTaskId, 950L)); // lag = 1000 - 950 = 50 → 
still hot
+        final StreamsGroupHeartbeatRequestData third = 
heartbeatState.buildRequestData();
+        assertEquals(950L, third.taskOffsets().get(0).offset());
     }
 
     @Test
@@ -926,10 +945,10 @@ class StreamsGroupHeartbeatRequestManagerTest {
         final StreamsRebalanceData.TaskId hotWarmup = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
         final StreamsRebalanceData.TaskId coldWarmup = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_2, 0);
 
-        final Map<StreamsRebalanceData.TaskId, Long> offsets = Map.of(
+        final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> offsets 
= new AtomicReference<>(Map.of(
             hotWarmup, 900L,   // lag = 1000 - 900 = 100 → hot
             coldWarmup, 500L   // lag = 1000 - 500 = 500 → cold
-        );
+        ));
         final Map<StreamsRebalanceData.TaskId, Long> endOffsets = Map.of(
             hotWarmup, 1000L,
             coldWarmup, 1000L
@@ -940,7 +959,7 @@ class StreamsGroupHeartbeatRequestManagerTest {
             Optional.of(RACK_ID),
             SUBTOPOLOGIES,
             CLIENT_TAGS,
-            () -> offsets,
+            offsets::get,
             () -> endOffsets
         );
         rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
@@ -966,9 +985,14 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertNotNull(first.taskOffsets());
         assertEquals(2, first.taskOffsets().size());
 
-        // Second HB without advancing time: assignmentChanged=false, 
taskOffsetIntervalPassed=false.
-        // The hot warmup makes hasAtLeastOneHotWarmupTask return true, 
triggering the send.
-        // Both warmups' offsets must appear in the resulting taskOffsets list.
+        // The hot warmup makes progress (the offset map changes). Without 
advancing time:
+        // assignmentChanged=false, taskOffsetIntervalPassed=false, but the 
hot warmup makes
+        // hasAtLeastOneHotWarmupTask return true, triggering the send. 
Because the map changed,
+        // the offsets for BOTH warmups must appear (the broker needs the 
complete picture).
+        offsets.set(Map.of(
+            hotWarmup, 950L,   // lag = 1000 - 950 = 50 → still hot
+            coldWarmup, 500L
+        ));
         final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
         assertNotNull(second.taskOffsets());
         assertEquals(2, second.taskOffsets().size());
@@ -978,7 +1002,7 @@ class StreamsGroupHeartbeatRequestManagerTest {
                 t -> new StreamsRebalanceData.TaskId(t.subtopologyId(), 
t.partition()),
                 StreamsGroupHeartbeatRequestData.TaskOffset::offset
             ));
-        assertEquals(900L, reportedOffsets.get(hotWarmup));
+        assertEquals(950L, reportedOffsets.get(hotWarmup));
         assertEquals(500L, reportedOffsets.get(coldWarmup));
     }
 
@@ -1092,6 +1116,174 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertNull(heartbeatState.buildRequestData().taskOffsets());
     }
 
+    @Test
+    public void testTaskOffsetsNotResentWhenUnchangedAcrossInterval() {
+        // The periodic task-offset interval trigger fires, but when neither 
the offsets nor the
+        // end-offsets changed since the last heartbeat, both fields are left 
null ("unchanged").
+        final StreamsRebalanceData.TaskId task = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final StreamsRebalanceData rebalanceData = 
newRebalanceDataWithStandbyOffsets(
+            task,
+            new AtomicReference<>(Map.of(task, 100L)),
+            new AtomicReference<>(Map.of(task, 200L))
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // First STABLE build: the assignment-changed trigger sends both 
fields.
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertNotNull(first.taskOffsets());
+        assertNotNull(first.taskEndOffsets());
+
+        // Advance past the interval. The interval trigger fires, but the 
values are unchanged.
+        time.sleep(1000);
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+        assertNull(second.taskEndOffsets());
+    }
+
+    @Test
+    public void testTaskOffsetsResentWhenChangedAcrossInterval() {
+        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);
+
+        assertEquals(100L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+
+        // The offset advanced; the next interval-triggered heartbeat resends 
it.
+        offsets.set(Map.of(task, 150L));
+        time.sleep(1000);
+        assertEquals(150L, 
heartbeatState.buildRequestData().taskOffsets().get(0).offset());
+    }
+
+    @Test
+    public void testTaskOffsetsAndEndOffsetsReportedIndependently() {
+        // A null field means "unchanged", and the two fields are independent: 
one may be sent
+        // while the other stays null.
+        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 AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> 
endOffsets =
+            new AtomicReference<>(Map.of(task, 200L));
+        final StreamsRebalanceData rebalanceData = 
newRebalanceDataWithStandbyOffsets(
+            task,
+            offsets,
+            endOffsets
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // First build sends both.
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertNotNull(first.taskOffsets());
+        assertNotNull(first.taskEndOffsets());
+
+        // Only the offsets change → only taskOffsets is sent; taskEndOffsets 
stays null.
+        offsets.set(Map.of(task, 120L));
+        time.sleep(1000);
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertEquals(120L, second.taskOffsets().get(0).offset());
+        assertNull(second.taskEndOffsets());
+
+        // Only the end-offsets change → only taskEndOffsets is sent; 
taskOffsets stays null.
+        endOffsets.set(Map.of(task, 220L));
+        time.sleep(1000);
+        final StreamsGroupHeartbeatRequestData third = 
heartbeatState.buildRequestData();
+        assertNull(third.taskOffsets());
+        assertEquals(220L, third.taskEndOffsets().get(0).offset());
+    }
+
+    @Test
+    public void testTaskOffsetsResentAfterReset() {
+        // reset() (called on every error/disconnect) clears the last-sent 
snapshot, so the next
+        // heartbeat resends the full offset state even if the values did not 
change. This is what
+        // makes "send only if changed" safe across coordinator failover 
(offsets are not persisted).
+        final StreamsRebalanceData.TaskId task = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final StreamsRebalanceData rebalanceData = 
newRebalanceDataWithStandbyOffsets(
+            task,
+            new AtomicReference<>(Map.of(task, 100L)),
+            new AtomicReference<>(Map.of(task, 200L))
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        assertNotNull(heartbeatState.buildRequestData().taskOffsets());
+
+        // Without a reset, the unchanged offsets would not be resent.
+        time.sleep(1000);
+        assertNull(heartbeatState.buildRequestData().taskOffsets());
+
+        // After a reset, the unchanged offsets are resent.
+        heartbeatState.reset();
+        final StreamsGroupHeartbeatRequestData afterReset = 
heartbeatState.buildRequestData();
+        assertEquals(100L, afterReset.taskOffsets().get(0).offset());
+        assertEquals(200L, afterReset.taskEndOffsets().get(0).offset());
+    }
+
+    @Test
+    public void 
testJoiningRecordsSentOffsetsSoFollowUpHeartbeatSkipsUnchanged() {
+        final StreamsRebalanceData.TaskId task = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final StreamsRebalanceData rebalanceData = 
newRebalanceDataWithStandbyOffsets(
+            task,
+            new AtomicReference<>(Map.of(task, 100L)),
+            new AtomicReference<>(Map.of(task, 200L))
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+
+        // Joining sends both fields and records them as last-sent.
+        when(membershipManager.state()).thenReturn(MemberState.JOINING);
+        final StreamsGroupHeartbeatRequestData joining = 
heartbeatState.buildRequestData();
+        assertNotNull(joining.taskOffsets());
+        assertNotNull(joining.taskEndOffsets());
+
+        // The immediately following non-joining heartbeat does not 
redundantly resend the
+        // unchanged offsets.
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+        final StreamsGroupHeartbeatRequestData followUp = 
heartbeatState.buildRequestData();
+        assertNull(followUp.taskOffsets());
+        assertNull(followUp.taskEndOffsets());
+    }
+
+    private StreamsRebalanceData newRebalanceDataWithStandbyOffsets(
+            final StreamsRebalanceData.TaskId standbyTaskId,
+            final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> 
taskOffsetSum,
+            final AtomicReference<Map<StreamsRebalanceData.TaskId, Long>> 
taskEndOffsetSum) {
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            taskOffsetSum::get,
+            taskEndOffsetSum::get
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(standbyTaskId),
+            Set.of(),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(100L);
+        return rebalanceData;
+    }
+
     private StreamsRebalanceData newRebalanceDataWithWarmup(final 
StreamsRebalanceData.TaskId warmupTaskId,
                                                             final long offset,
                                                             final long 
endOffset,

Reply via email to