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 d9b215da915 MINOR: handle invalid StreamsHeartbeatRequest (#22728)
d9b215da915 is described below
commit d9b215da9156a5ce080cea30c0eca73851c9d5e8
Author: Matthias J. Sax <[email protected]>
AuthorDate: Thu Jul 2 23:32:33 2026 -0700
MINOR: handle invalid StreamsHeartbeatRequest (#22728)
A StreamsHeartbeatRequest must not contain duplicate task-offset
entries, and should be rejected as invalid request. We also need to
avoid to mutate in-memory state for an invalid request.
Reviewers: David Jacot <[email protected]>
---
.../coordinator/group/GroupMetadataManager.java | 20 ++--
.../group/streams/MemberTaskOffsets.java | 19 +++-
.../group/GroupMetadataManagerTest.java | 106 +++++++++++++++++++++
.../group/streams/MemberTaskOffsetsTest.java | 19 ++++
4 files changed, 151 insertions(+), 13 deletions(-)
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
index 9643099e088..38fee045c3d 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java
@@ -2116,14 +2116,6 @@ public class GroupMetadataManager {
);
}
- // Store the latest task changelog offsets/end-offsets reported by the
member. These are transient telemetry
- // (used by the assignor to estimate task lag) and are not persisted.
Task offsets and end-offsets are reported
- // independently: a null list means "unchanged since the last
heartbeat", so we retain the previously reported
- // value for whichever of the two is null and only update when at
least one is reported.
- if (taskOffsets != null || taskEndOffsets != null) {
- group.updateTaskOffsets(memberId,
group.taskOffsets(memberId).update(taskOffsets, taskEndOffsets));
- }
-
// 1. Create or update the member.
StreamsGroupMember.Builder updatedMemberBuilder = new
StreamsGroupMember.Builder(member)
.maybeUpdateInstanceId(Optional.ofNullable(instanceId))
@@ -2194,6 +2186,18 @@ public class GroupMetadataManager {
throwIfRequestContainsInvalidTasks(subtopologySortedMap,
ownedStandbyTasks);
throwIfRequestContainsInvalidTasks(subtopologySortedMap,
ownedWarmupTasks);
}
+
+ // Store the latest task changelog offsets/end-offsets reported by the
member. These are transient telemetry
+ // (used by the assignor to estimate task lag) and are not persisted.
Task offsets and end-offsets are reported
+ // independently: a null list means "unchanged since the last
heartbeat", so we retain the previously reported
+ // value for whichever of the two is null and only update when at
least one is reported.
+ // This must run after the task validation above: it mutates a
non-timeline map that the coordinator runtime
+ // does not roll back, so updating it before validation could leave an
orphaned entry for a member whose
+ // (joining) heartbeat is then rejected.
+ if (taskOffsets != null || taskEndOffsets != null) {
+ group.updateTaskOffsets(memberId,
group.taskOffsets(memberId).update(taskOffsets, taskEndOffsets));
+ }
+
// We validated a topology that was not validated before, so bump the
group epoch as we may have to reassign tasks.
if (validatedTopologyEpoch != group.validatedTopologyEpoch()) {
bumpGroupEpoch = true;
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
index 40c863b4b20..ed2434a79bf 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsets.java
@@ -16,12 +16,13 @@
*/
package org.apache.kafka.coordinator.group.streams;
+import org.apache.kafka.common.errors.InvalidRequestException;
import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData;
import org.apache.kafka.coordinator.group.streams.assignor.TaskId;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.stream.Collectors;
/**
* The latest per-task cumulative changelog offsets and end-offsets that a
member reported in its heartbeat.
@@ -57,9 +58,17 @@ public record MemberTaskOffsets(Map<TaskId, Long>
taskOffsets, Map<TaskId, Long>
}
private static Map<TaskId, Long> toTaskIdMap(final
List<StreamsGroupHeartbeatRequestData.TaskOffset> taskOffsets) {
- return taskOffsets.stream().collect(Collectors.toMap(
- taskOffset -> new TaskId(taskOffset.subtopologyId(),
taskOffset.partition()),
- StreamsGroupHeartbeatRequestData.TaskOffset::offset
- ));
+ final Map<TaskId, Long> result = new HashMap<>(taskOffsets.size());
+ for (final StreamsGroupHeartbeatRequestData.TaskOffset taskOffset :
taskOffsets) {
+ final TaskId taskId = new TaskId(taskOffset.subtopologyId(),
taskOffset.partition());
+ // The reported values come straight from the client heartbeat and
the protocol does not enforce uniqueness
+ // of (subtopologyId, partition). Reject a duplicate with a clear
client error rather than silently picking
+ // one of the values.
+ if (result.putIfAbsent(taskId, taskOffset.offset()) != null) {
+ throw new InvalidRequestException("Task offsets contain a
duplicate entry for subtopology "
+ + taskOffset.subtopologyId() + " and partition " +
taskOffset.partition() + ".");
+ }
+ }
+ return result;
}
}
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
index 9f348600f4c..f3568dc089b 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/GroupMetadataManagerTest.java
@@ -134,6 +134,7 @@ import
org.apache.kafka.coordinator.group.modern.share.ShareGroup.InitMapValue;
import org.apache.kafka.coordinator.group.modern.share.ShareGroupBuilder;
import org.apache.kafka.coordinator.group.modern.share.ShareGroupConfig;
import org.apache.kafka.coordinator.group.modern.share.ShareGroupMember;
+import org.apache.kafka.coordinator.group.streams.MemberTaskOffsets;
import org.apache.kafka.coordinator.group.streams.MockTaskAssignor;
import
org.apache.kafka.coordinator.group.streams.StreamsCoordinatorRecordHelpers;
import org.apache.kafka.coordinator.group.streams.StreamsGroup;
@@ -18576,6 +18577,111 @@ public class GroupMetadataManagerTest {
assertEquals("Task 3 for subtopology subtopology1 is invalid. Number
of tasks for this subtopology: 3", e2.getMessage());
}
+ @Test
+ public void
testStreamsGroupHeartbeatDoesNotStoreTaskOffsetsWhenOwnedTasksAreInvalid() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 3)
+ .buildCoordinatorMetadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, 10)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+
.setAssignedTasks(mkTasksTupleWithCommonEpoch(TaskRole.ACTIVE, 10,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+ .build())
+ .withTopology(StreamsTopology.fromHeartbeatRequest(topology))
+ .withTargetAssignment(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+ .withTargetAssignmentEpoch(10)
+ )
+ .build();
+
+ // The heartbeat reports task offsets but also claims an invalid owned
task, so it is rejected during
+ // validation. The reported offsets must not be stored: updating the
(non-timeline) task-offsets map before
+ // validation would leave an orphaned entry that the coordinator
runtime does not roll back.
+ assertThrows(InvalidRequestException.class, () ->
context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(10)
+ .setActiveTasks(List.of(new
StreamsGroupHeartbeatRequestData.TaskIds()
+ .setSubtopologyId(subtopology1)
+ .setPartitions(List.of(3)))) // partition 3 is out of
range (only 0, 1, 2 exist)
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of())
+ .setTaskOffsets(List.of(new
StreamsGroupHeartbeatRequestData.TaskOffset()
+
.setSubtopologyId(subtopology1).setPartition(0).setOffset(10L)))
+ .setTaskEndOffsets(List.of(new
StreamsGroupHeartbeatRequestData.TaskOffset()
+
.setSubtopologyId(subtopology1).setPartition(0).setOffset(20L)))));
+
+ StreamsGroup group =
context.groupMetadataManager.streamsGroup(groupId);
+ assertEquals(MemberTaskOffsets.EMPTY, group.taskOffsets(memberId));
+ }
+
+ @Test
+ public void testStreamsGroupHeartbeatRejectsDuplicateTaskOffsets() {
+ String groupId = "fooup";
+ String memberId = Uuid.randomUuid().toString();
+ String subtopology1 = "subtopology1";
+ String fooTopicName = "foo";
+ Uuid fooTopicId = Uuid.randomUuid();
+ Topology topology = new Topology().setSubtopologies(List.of(
+ new
Subtopology().setSubtopologyId(subtopology1).setSourceTopics(List.of(fooTopicName))
+ ));
+
+ MockTaskAssignor assignor = new MockTaskAssignor("sticky");
+ GroupMetadataManagerTestContext context = new
GroupMetadataManagerTestContext.Builder()
+ .withStreamsGroupTaskAssignors(List.of(assignor))
+ .withMetadataImage(new MetadataImageBuilder()
+ .addTopic(fooTopicId, fooTopicName, 3)
+ .buildCoordinatorMetadataImage())
+ .withStreamsGroup(new StreamsGroupBuilder(groupId, 10)
+ .withMember(streamsGroupMemberBuilderWithDefaults(memberId)
+ .setMemberEpoch(10)
+ .setPreviousMemberEpoch(10)
+
.setAssignedTasks(mkTasksTupleWithCommonEpoch(TaskRole.ACTIVE, 10,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+ .build())
+ .withTopology(StreamsTopology.fromHeartbeatRequest(topology))
+ .withTargetAssignment(memberId,
TaskAssignmentTestUtil.mkTasksTuple(TaskRole.ACTIVE,
+ TaskAssignmentTestUtil.mkTasks(subtopology1, 0, 1, 2)))
+ .withTargetAssignmentEpoch(10)
+ )
+ .build();
+
+ // The owned tasks are valid, but the reported task offsets contain a
duplicate (subtopologyId, partition)
+ // entry. This is a malformed request and must be rejected with a
clear client error; nothing is stored.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () -> context.streamsGroupHeartbeat(
+ new StreamsGroupHeartbeatRequestData()
+ .setGroupId(groupId)
+ .setMemberId(memberId)
+ .setMemberEpoch(10)
+ .setActiveTasks(List.of(new
StreamsGroupHeartbeatRequestData.TaskIds()
+ .setSubtopologyId(subtopology1)
+ .setPartitions(List.of(0, 1, 2))))
+ .setStandbyTasks(List.of())
+ .setWarmupTasks(List.of())
+ .setTaskOffsets(List.of(
+ new
StreamsGroupHeartbeatRequestData.TaskOffset().setSubtopologyId(subtopology1).setPartition(0).setOffset(10L),
+ new
StreamsGroupHeartbeatRequestData.TaskOffset().setSubtopologyId(subtopology1).setPartition(0).setOffset(11L)))));
+ assertEquals("Task offsets contain a duplicate entry for subtopology "
+ subtopology1 + " and partition 0.", e.getMessage());
+
+ StreamsGroup group =
context.groupMetadataManager.streamsGroup(groupId);
+ assertEquals(MemberTaskOffsets.EMPTY, group.taskOffsets(memberId));
+ }
+
@Test
public void testStreamsGroupHeartbeatStoresTaskOffsetsWithoutPersisting() {
String groupId = "fooup";
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
index 68ff6b0e370..941c5c3cc97 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/MemberTaskOffsetsTest.java
@@ -16,6 +16,7 @@
*/
package org.apache.kafka.coordinator.group.streams;
+import org.apache.kafka.common.errors.InvalidRequestException;
import org.apache.kafka.common.message.StreamsGroupHeartbeatRequestData;
import org.apache.kafka.coordinator.group.streams.assignor.TaskId;
@@ -25,6 +26,7 @@ import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
public class MemberTaskOffsetsTest {
@@ -54,6 +56,23 @@ public class MemberTaskOffsetsTest {
);
}
+ @Test
+ public void shouldRejectDuplicateTaskEntries() {
+ // The protocol does not enforce uniqueness of (subtopologyId,
partition) within a reported list. A duplicate
+ // entry is a malformed request and must be rejected with a clear
client error rather than silently resolved.
+ InvalidRequestException e =
assertThrows(InvalidRequestException.class, () ->
MemberTaskOffsets.EMPTY.update(
+ List.of(taskOffset("sub-1", 0, 10L), taskOffset("sub-1", 0, 42L)),
+ null
+ ));
+ assertEquals("Task offsets contain a duplicate entry for subtopology
sub-1 and partition 0.", e.getMessage());
+
+ // The check also applies to the end-offsets list.
+ assertThrows(InvalidRequestException.class, () ->
MemberTaskOffsets.EMPTY.update(
+ null,
+ List.of(taskOffset("sub-1", 0, 15L), taskOffset("sub-1", 0, 55L))
+ ));
+ }
+
@Test
public void shouldRetainBothMapsWhenBothListsAreNull() {
MemberTaskOffsets previous = new MemberTaskOffsets(