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 844de98644a KAFKA-20116: Send task-(end)-offset to broker (3/N) 
(#21803)
844de98644a is described below

commit 844de98644a522aefb666dc216cec2e71336363f
Author: Matthias J. Sax <[email protected]>
AuthorDate: Thu Jun 25 11:37:55 2026 -0700

    KAFKA-20116: Send task-(end)-offset to broker (3/N) (#21803)
    
    This PR add code to send taskOffsetSum and taskEndOffsetSum
    via StreamsHeartbeatRequest to the broker.
    
    Part of KIP-1071.
    
    Reviewers: Lucas Brutschy <[email protected]>
---
 .../StreamsGroupHeartbeatRequestManager.java       |  75 ++++-
 .../requests/StreamsGroupHeartbeatRequest.java     |   5 +
 .../StreamsGroupHeartbeatRequestManagerTest.java   | 356 ++++++++++++++++++++-
 .../requests/StreamsGroupHeartbeatRequestTest.java |  65 ++++
 .../coordinator/group/GroupCoordinatorService.java |   2 -
 .../group/GroupCoordinatorServiceTest.java         |  34 --
 6 files changed, 485 insertions(+), 52 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 6ae3e048c7c..d59a2f2c90f 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
@@ -90,20 +90,24 @@ public class StreamsGroupHeartbeatRequestManager implements 
RequestManager {
         private final StreamsMembershipManager membershipManager;
         private final int rebalanceTimeoutMs;
         private final StreamsRebalanceData streamsRebalanceData;
+        private final Time time;
         private final LastSentFields lastSentFields = new LastSentFields();
         private int endpointInformationEpoch = -1;
-
+        private long lastTaskOffsetIntervalTs = -1;
 
         public HeartbeatState(final StreamsRebalanceData streamsRebalanceData,
                               final StreamsMembershipManager membershipManager,
-                              final int rebalanceTimeoutMs) {
+                              final int rebalanceTimeoutMs,
+                              final Time time) {
             this.membershipManager = membershipManager;
             this.streamsRebalanceData = streamsRebalanceData;
             this.rebalanceTimeoutMs = rebalanceTimeoutMs;
+            this.time = time;
         }
 
         public void reset() {
             lastSentFields.reset();
+            lastTaskOffsetIntervalTs = -1L;
         }
 
         public int endpointInformationEpoch() {
@@ -147,19 +151,80 @@ 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()));
+                lastTaskOffsetIntervalTs = time.milliseconds();
             } else {
-                StreamsRebalanceData.Assignment reconciledAssignment = 
streamsRebalanceData.reconciledAssignment();
-                if (!reconciledAssignment.equals(lastSentFields.assignment)) {
+                final StreamsRebalanceData.Assignment reconciledAssignment = 
streamsRebalanceData.reconciledAssignment();
+                final boolean assignmentChanged = 
!reconciledAssignment.equals(lastSentFields.assignment);
+
+                if (assignmentChanged) {
                     
data.setActiveTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.activeTasks()));
                     
data.setStandbyTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.standbyTasks()));
                     
data.setWarmupTasks(fromStreamsToHeartbeatRequest(reconciledAssignment.warmupTasks()));
                     lastSentFields.assignment = reconciledAssignment;
                 }
+
+                // call both method only once, as they invoke an expensive 
`supplier`
+                final Map<StreamsRebalanceData.TaskId, Long> taskOffsetSum = 
streamsRebalanceData.taskOffsetSum();
+                final Map<StreamsRebalanceData.TaskId, Long> taskEndOffsetSum 
= streamsRebalanceData.taskEndOffsetSum();
+
+                if (assignmentChanged
+                    || taskOffsetIntervalPassed()
+                    || 
hasAtLeastOneHotWarmupTask(reconciledAssignment.warmupTasks(), taskOffsetSum, 
taskEndOffsetSum)
+                ) {
+
+                    // TODO: send only if changed this last time
+                    data.setTaskOffsets(convertToList(taskOffsetSum));
+                    data.setTaskEndOffsets(convertToList(taskEndOffsetSum));
+
+                    lastTaskOffsetIntervalTs = time.milliseconds();
+                }
             }
             
data.setShutdownApplication(streamsRebalanceData.shutdownRequested());
             return data;
         }
 
+        private static List<StreamsGroupHeartbeatRequestData.TaskOffset> 
convertToList(Map<StreamsRebalanceData.TaskId, Long> offsetsMap) {
+            return offsetsMap.entrySet().stream().map(
+                    entry -> new StreamsGroupHeartbeatRequestData.TaskOffset()
+                        .setSubtopologyId(entry.getKey().subtopologyId())
+                        .setPartition(entry.getKey().partitionId())
+                        .setOffset(entry.getValue()))
+                .collect(Collectors.toList());
+        }
+
+        private boolean taskOffsetIntervalPassed() {
+            return lastTaskOffsetIntervalTs + 
streamsRebalanceData.taskOffsetIntervalMs() <= time.milliseconds();
+        }
+
+        private boolean hasAtLeastOneHotWarmupTask(
+            final Set<StreamsRebalanceData.TaskId> warmupTasks,
+            final Map<StreamsRebalanceData.TaskId, Long> taskOffsetSum,
+            final Map<StreamsRebalanceData.TaskId, Long> taskEndOffsetSum
+        ) {
+            if (warmupTasks.isEmpty()) {
+                return false;
+            }
+
+            final long acceptableRecoveryLag = 
streamsRebalanceData.acceptableRecoveryLag();
+
+            return warmupTasks.stream()
+                .anyMatch(taskId -> {
+                    final Long offset = taskOffsetSum.get(taskId);
+                    final Long endOffset = taskEndOffsetSum.get(taskId);
+
+                    // offset and endOffset might not be known,
+                    // or be capped at MAX_VALUE due to overflow
+                    if (offset == null || offset == Long.MAX_VALUE
+                        || endOffset == null || endOffset == Long.MAX_VALUE) {
+                        return false;
+                    }
+
+                    return endOffset - offset <= acceptableRecoveryLag;
+                });
+        }
+
         private static List<StreamsGroupHeartbeatRequestData.TaskIds> 
fromStreamsToHeartbeatRequest(final Set<StreamsRebalanceData.TaskId> tasks) {
             return tasks.stream()
                 .collect(
@@ -327,7 +392,7 @@ public class StreamsGroupHeartbeatRequestManager implements 
RequestManager {
         this.maxPollIntervalMs = 
config.getInt(CommonClientConfigs.MAX_POLL_INTERVAL_MS_CONFIG);
         long retryBackoffMs = 
config.getLong(ConsumerConfig.RETRY_BACKOFF_MS_CONFIG);
         long retryBackoffMaxMs = 
config.getLong(ConsumerConfig.RETRY_BACKOFF_MAX_MS_CONFIG);
-        this.heartbeatState = new HeartbeatState(streamsRebalanceData, 
membershipManager, maxPollIntervalMs);
+        this.heartbeatState = new HeartbeatState(streamsRebalanceData, 
membershipManager, maxPollIntervalMs, time);
         this.heartbeatRequestState = new HeartbeatRequestState(
             logContext,
             time,
diff --git 
a/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequest.java
 
b/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequest.java
index 26d1592f040..74d23a0b704 100644
--- 
a/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequest.java
+++ 
b/clients/src/main/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequest.java
@@ -45,6 +45,11 @@ public class StreamsGroupHeartbeatRequest extends 
AbstractRequest {
 
         @Override
         public StreamsGroupHeartbeatRequest build(short version) {
+            if (version == 0) {
+                // TaskOffsets/TaskEndOffsets are only supported by brokers 
supporting v1+ request versions
+                data.setTaskOffsets(null);
+                data.setTaskEndOffsets(null);
+            }
             return new StreamsGroupHeartbeatRequest(data, version);
         }
 
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 c2083fbafd0..315c58c0b89 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
@@ -74,6 +74,7 @@ import static 
org.apache.kafka.common.requests.StreamsGroupHeartbeatRequest.LEAV
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -654,7 +655,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1000
+                1000,
+                time
             );
 
         StreamsGroupHeartbeatRequestData requestData1 = 
heartbeatState.buildRequestData();
@@ -687,7 +689,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -709,7 +712,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -737,7 +741,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1000
+                1000,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -795,6 +800,328 @@ class StreamsGroupHeartbeatRequestManagerTest {
         assertNull(nonJoiningRequestData.topology());
     }
 
+    @Test
+    public void testHotWarmupTaskDisabledWhenAcceptableRecoveryLagIsNegative() 
{
+        // A v0 broker (or a yet-unreachable coordinator) leaves 
acceptableRecoveryLag at its
+        // -1 default. In that case `hasHotWarmupTask` must always return 
`false`
+        //
+        // Note: for v0 broker case, we would actually never expect to get 
warmup task assigned,
+        // so the test code below does not fully mimic reality, but rather 
test a corner case
+        // which should never hit in reality
+        final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final Map<StreamsRebalanceData.TaskId, Long> warmupOffsets = Map.of(
+            new StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0), 42L
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            () -> warmupOffsets,
+            Map::of
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(),
+            Set.of(warmupTaskId), // this is technically incorrect, but 
ensures that `acceptable.recovery.lag == -1` is tested correctly
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(-1L);
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // first HB always sends the offset
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertEquals(42L, 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 
acceptableRecoveryLag == -1 it must return false.
+        // The result is that no TaskOffsets field is set on the request.
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+    }
+
+    @Test
+    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.
+        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 StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // first HB always sends the offset
+        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
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertEquals(900L, second.taskOffsets().get(0).offset());
+    }
+
+    @Test
+    public void testHotWarmupTaskDisabledWhenLagAboveThreshold() {
+        // Warmup whose lag exceeds acceptableRecoveryLag must NOT trigger an 
early send.
+        final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final StreamsRebalanceData rebalanceData = newRebalanceDataWithWarmup(
+            warmupTaskId,
+            500L,  // offset
+            1000L, // endOffset → lag = 500
+            100L   // acceptableRecoveryLag (lag 500 > 100)
+        );
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // first HB always sends the offset
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertEquals(500L, first.taskOffsets().get(0).offset());
+
+        // high lag -- don't send
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+    }
+
+    @Test
+    public void testTaskOffsetsForAllWarmupsAreReportedWhenAtLeastOneIsHot() {
+        // Two warmup tasks: the first is hot (lag below threshold), the 
second is cold
+        // (lag above threshold). hasAtLeastOneHotWarmupTask returns true 
because the first
+        // is hot — and the resulting heartbeat must carry the taskOffsets for 
BOTH warmups,
+        // not just the hot one. The broker needs the complete picture to 
drive its own
+        // lag-based promotion logic across the whole assignment.
+        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(
+            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
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            () -> offsets,
+            () -> endOffsets
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(),
+            Set.of(hotWarmup, coldWarmup),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(100L);
+
+        final StreamsGroupHeartbeatRequestManager.HeartbeatState 
heartbeatState =
+            new StreamsGroupHeartbeatRequestManager.HeartbeatState(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // First HB: assignment-changed trigger always sends offsets for both 
warmups.
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        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.
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNotNull(second.taskOffsets());
+        assertEquals(2, second.taskOffsets().size());
+
+        final Map<StreamsRebalanceData.TaskId, Long> reportedOffsets = 
second.taskOffsets().stream()
+            .collect(Collectors.toMap(
+                t -> new StreamsRebalanceData.TaskId(t.subtopologyId(), 
t.partition()),
+                StreamsGroupHeartbeatRequestData.TaskOffset::offset
+            ));
+        assertEquals(900L, reportedOffsets.get(hotWarmup));
+        assertEquals(500L, reportedOffsets.get(coldWarmup));
+    }
+
+    @Test
+    public void testHotWarmupTaskDisabledWhenEndOffsetMissing() {
+        // Without an end-offset entry for the warmup task, lag cannot be 
computed.
+        // hasHotWarmupTask must return false (safe fallback) — the broker 
will still
+        // receive whatever partial information arrives via the normal 
interval-driven path.
+        final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final Map<StreamsRebalanceData.TaskId, Long> offsets = Map.of(
+            new StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0), 900L
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            () -> offsets,
+            Map::of // no end-offsets
+        );
+        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(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // first HB always sends the offset
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertEquals(900L, first.taskOffsets().get(0).offset());
+
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+    }
+
+    @Test
+    public void testHotWarmupTaskDisabledWhenOffsetMissing() {
+        // Symmetric to the end-offset-missing case: if the warmup has an 
end-offset entry but
+        // its offset entry is absent from `taskOffsetSum`, lag cannot be 
computed.
+        final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+        final Map<StreamsRebalanceData.TaskId, Long> endOffsets = Map.of(
+            new StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0), 1000L
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            Map::of, // no offsets
+            () -> endOffsets
+        );
+        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(
+                rebalanceData,
+                membershipManager,
+                1234,
+                time
+            );
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+
+        // first HB always sends the offset
+        final StreamsGroupHeartbeatRequestData first = 
heartbeatState.buildRequestData();
+        assertNotNull(first.taskOffsets());
+
+        final StreamsGroupHeartbeatRequestData second = 
heartbeatState.buildRequestData();
+        assertNull(second.taskOffsets());
+    }
+
+    @Test
+    public void testHotWarmupTaskDisabledWhenOffsetOrEndOffsetOverflowed() {
+        // Either side may be pinned to Long.MAX_VALUE on overflow during 
cross-store summing
+        // (see StreamsPartitionAssignor.computeEndOffsetSumsByTask). The 
arithmetic
+        // endOffset - offset would otherwise produce a misleading result, so 
the predicate
+        // must conservatively return false.
+        final StreamsRebalanceData.TaskId warmupTaskId = new 
StreamsRebalanceData.TaskId(SUBTOPOLOGY_NAME_1, 0);
+
+        // offset = MAX_VALUE → not hot.
+        StreamsRebalanceData rebalanceData = 
newRebalanceDataWithWarmup(warmupTaskId, Long.MAX_VALUE, 1000L, 100L);
+        StreamsGroupHeartbeatRequestManager.HeartbeatState heartbeatState =
+            new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        when(membershipManager.state()).thenReturn(MemberState.STABLE);
+        assertNotNull(heartbeatState.buildRequestData().taskOffsets()); // 
first call: assignmentChanged trigger
+        assertNull(heartbeatState.buildRequestData().taskOffsets());    // 
second: hasAtLeastOneHotWarmupTask must bail
+
+        // endOffset = MAX_VALUE → not hot.
+        rebalanceData = newRebalanceDataWithWarmup(warmupTaskId, 900L, 
Long.MAX_VALUE, 100L);
+        heartbeatState = new 
StreamsGroupHeartbeatRequestManager.HeartbeatState(rebalanceData, 
membershipManager, 1234, time);
+        assertNotNull(heartbeatState.buildRequestData().taskOffsets());
+        assertNull(heartbeatState.buildRequestData().taskOffsets());
+    }
+
+    private StreamsRebalanceData newRebalanceDataWithWarmup(final 
StreamsRebalanceData.TaskId warmupTaskId,
+                                                            final long offset,
+                                                            final long 
endOffset,
+                                                            final long 
acceptableRecoveryLag) {
+        final Map<StreamsRebalanceData.TaskId, Long> offsets = Map.of(
+            new StreamsRebalanceData.TaskId(warmupTaskId.subtopologyId(), 
warmupTaskId.partitionId()), offset
+        );
+        final Map<StreamsRebalanceData.TaskId, Long> endOffsets = Map.of(
+                new StreamsRebalanceData.TaskId(warmupTaskId.subtopologyId(), 
warmupTaskId.partitionId()), endOffset
+        );
+        final StreamsRebalanceData rebalanceData = new StreamsRebalanceData(
+            PROCESS_ID,
+            Optional.of(ENDPOINT),
+            Optional.of(RACK_ID),
+            SUBTOPOLOGIES,
+            CLIENT_TAGS,
+            () -> offsets,
+            () -> endOffsets
+        );
+        rebalanceData.setReconciledAssignment(new 
StreamsRebalanceData.Assignment(
+            Set.of(),
+            Set.of(),
+            Set.of(warmupTaskId),
+            true
+        ));
+        rebalanceData.setTaskOffsetIntervalMs(1000);
+        rebalanceData.setAcceptableRecoveryLag(acceptableRecoveryLag);
+        return rebalanceData;
+    }
+
     private <V> boolean isSorted(List<V> collection, Comparator<V> comparator) 
{
         for (int i = 1; i < collection.size(); i++) {
             if (comparator.compare(collection.get(i - 1), collection.get(i)) > 
0) {
@@ -812,7 +1139,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                rebalanceTimeoutMs
+                rebalanceTimeoutMs,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -834,7 +1162,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -856,7 +1185,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -879,7 +1209,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -902,7 +1233,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(MemberState.JOINING);
 
@@ -1014,7 +1346,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(memberState);
         streamsRebalanceData.setReconciledAssignment(
@@ -1076,7 +1409,8 @@ class StreamsGroupHeartbeatRequestManagerTest {
             new StreamsGroupHeartbeatRequestManager.HeartbeatState(
                 streamsRebalanceData,
                 membershipManager,
-                1234
+                1234,
+                time
             );
         when(membershipManager.state()).thenReturn(memberState);
 
diff --git 
a/clients/src/test/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequestTest.java
 
b/clients/src/test/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequestTest.java
new file mode 100644
index 00000000000..a480e58e227
--- /dev/null
+++ 
b/clients/src/test/java/org/apache/kafka/common/requests/StreamsGroupHeartbeatRequestTest.java
@@ -0,0 +1,65 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.common.requests;
+
+import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+public class StreamsGroupHeartbeatRequestTest {
+
+    private static StreamsGroupHeartbeatRequestData dataWithTaskOffsets() {
+        return new StreamsGroupHeartbeatRequestData()
+            .setGroupId("group")
+            .setMemberId("member")
+            .setTaskOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()
+                .setSubtopologyId("sub")
+                .setPartition(0)
+                .setOffset(42L)))
+            .setTaskEndOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()
+                .setSubtopologyId("sub")
+                .setPartition(0)
+                .setOffset(100L)));
+    }
+
+    @Test
+    public void testBuildClearsTaskOffsetsForVersion0() {
+        // TaskOffsets/TaskEndOffsets are only consumed by brokers supporting 
v1+.
+        // Against a v0 coordinator they must be omitted.
+        StreamsGroupHeartbeatRequest request = new 
StreamsGroupHeartbeatRequest.Builder(dataWithTaskOffsets())
+            .build((short) 0);
+
+        assertNull(request.data().taskOffsets());
+        assertNull(request.data().taskEndOffsets());
+    }
+
+    @Test
+    public void testBuildKeepsTaskOffsetsForVersion1() {
+        StreamsGroupHeartbeatRequest request = new 
StreamsGroupHeartbeatRequest.Builder(dataWithTaskOffsets())
+            .build((short) 1);
+
+        assertEquals(1, request.data().taskOffsets().size());
+        assertEquals(42L, request.data().taskOffsets().get(0).offset());
+        assertEquals(1, request.data().taskEndOffsets().size());
+        assertEquals(100L, request.data().taskEndOffsets().get(0).offset());
+    }
+}
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java
index dfe017f8108..8fc5719e76a 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupCoordinatorService.java
@@ -617,8 +617,6 @@ public class GroupCoordinatorService implements 
GroupCoordinator {
     private static void 
throwIfStreamsGroupHeartbeatRequestIsUsingUnsupportedFeatures(
         StreamsGroupHeartbeatRequestData request
     ) throws InvalidRequestException {
-        throwIfNotNull(request.taskOffsets(), "TaskOffsets are not supported 
yet.");
-        throwIfNotNull(request.taskEndOffsets(), "TaskEndOffsets are not 
supported yet.");
         throwIfNotNullOrEmpty(request.warmupTasks(), "WarmupTasks are not 
supported yet.");
         if (request.topology() != null) {
             for (StreamsGroupHeartbeatRequestData.Subtopology subtopology : 
request.topology().subtopologies()) {
diff --git 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java
 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java
index 16adc5b0618..458d150b01f 100644
--- 
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java
+++ 
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupCoordinatorServiceTest.java
@@ -597,40 +597,6 @@ public class GroupCoordinatorServiceTest {
         AuthorizableRequestContext context = 
mock(AuthorizableRequestContext.class);
         when(context.requestVersion()).thenReturn((int) 
ApiKeys.STREAMS_GROUP_HEARTBEAT.latestVersion());
 
-        assertEquals(
-            new StreamsGroupHeartbeatResult(
-                new StreamsGroupHeartbeatResponseData()
-                    .setErrorCode(Errors.INVALID_REQUEST.code())
-                    .setErrorMessage("TaskOffsets are not supported yet."),
-                Map.of(),
-                -1,
-                -1,
-                -1
-            ),
-            service.streamsGroupHeartbeat(
-                context,
-                new StreamsGroupHeartbeatRequestData()
-                    .setTaskOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()))
-            ).get(5, TimeUnit.SECONDS)
-        );
-
-        assertEquals(
-            new StreamsGroupHeartbeatResult(
-                new StreamsGroupHeartbeatResponseData()
-                    .setErrorCode(Errors.INVALID_REQUEST.code())
-                    .setErrorMessage("TaskEndOffsets are not supported yet."),
-                Map.of(),
-                -1,
-                -1,
-                -1
-            ),
-            service.streamsGroupHeartbeat(
-                context,
-                new StreamsGroupHeartbeatRequestData()
-                    .setTaskEndOffsets(List.of(new 
StreamsGroupHeartbeatRequestData.TaskOffset()))
-            ).get(5, TimeUnit.SECONDS)
-        );
-
         assertEquals(
             new StreamsGroupHeartbeatResult(
                 new StreamsGroupHeartbeatResponseData()


Reply via email to