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,