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