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 89ccd6a1260 KAFKA-20868: Add warmup tasks to IQ metadata (#23006)
89ccd6a1260 is described below
commit 89ccd6a1260130e1a5ea83a0aa98a7d397482ace
Author: gabriellafu <[email protected]>
AuthorDate: Tue Aug 4 21:10:24 2026 -0400
KAFKA-20868: Add warmup tasks to IQ metadata (#23006)
This PR adds warmupTopicPartitions to standbyTopicPartitions, to be
reported to the Kafka Streams clients as part of "host endpoint
metadata"via StreamsHeartbeatResponse, to allow querying warmup
tasks via IQ.
Part of KIP-1071.
Reviewers: Matthias J. Sax <[email protected]>
---
.../message/StreamsGroupHeartbeatResponse.json | 2 +-
.../group/streams/TasksTupleWithEpochs.java | 7 +++++
.../topics/EndpointToPartitionsManager.java | 5 ++--
.../topics/EndpointToPartitionsManagerTest.java | 30 ++++++++++++++++++++++
4 files changed, 41 insertions(+), 3 deletions(-)
diff --git
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
index a63ddbb4af2..2e6bb8a7706 100644
---
a/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
+++
b/clients/src/main/resources/common/message/StreamsGroupHeartbeatResponse.json
@@ -83,7 +83,7 @@
{ "name": "ActivePartitions", "type": "[]TopicPartition", "versions":
"0+",
"about": "All topic partitions materialized by active tasks on the
node" },
{ "name": "StandbyPartitions", "type": "[]TopicPartition", "versions":
"0+",
- "about": "All topic partitions materialized by standby tasks on the
node" }
+ "about": "All topic partitions materialized by standby and warm-up
tasks on the node" }
]
}
],
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TasksTupleWithEpochs.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TasksTupleWithEpochs.java
index 99220d466bc..f48e788730c 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TasksTupleWithEpochs.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/TasksTupleWithEpochs.java
@@ -68,6 +68,13 @@ public record TasksTupleWithEpochs(Map<String, Map<Integer,
Integer>> activeTask
return activeTasksWithEpochs.isEmpty() && standbyTasks.isEmpty() &&
warmupTasks.isEmpty();
}
+ /**
+ * @return standby and warm-up tasks merged, since warm-up tasks are
executed as standby tasks on the client.
+ */
+ public Map<String, Set<Integer>> standbyAndWarmupTasks() {
+ return mergeTasks(standbyTasks, warmupTasks);
+ }
+
/**
* Merges this task tuple with another task tuple.
* For overlapping active tasks, epochs from the other tuple take
precedence.
diff --git
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManager.java
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManager.java
index 0798eddf45c..5732fdf2760 100644
---
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManager.java
+++
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManager.java
@@ -72,11 +72,12 @@ public class EndpointToPartitionsManager {
Map.Entry::getKey,
entry -> entry.getValue().keySet()
));
- Map<String, Set<Integer>> standbyTasks =
streamsGroupMember.assignedTasks().standbyTasks();
+ // Warm-up tasks are executed as standby tasks on the client, so they
are reported as standby partitions, as in the classic protocol.
+ Map<String, Set<Integer>> standbyAndWarmupTasks =
streamsGroupMember.assignedTasks().standbyAndWarmupTasks();
endpointToPartitions.setUserEndpoint(responseEndpoint);
Map<String, ConfiguredSubtopology> configuredSubtopologies =
streamsGroup.configuredTopology().flatMap(ConfiguredTopology::subtopologies).get();
List<StreamsGroupHeartbeatResponseData.TopicPartition>
activeTopicPartitions = topicPartitions(activeTasks, configuredSubtopologies,
metadataImage);
- List<StreamsGroupHeartbeatResponseData.TopicPartition>
standbyTopicPartitions = topicPartitions(standbyTasks, configuredSubtopologies,
metadataImage);
+ List<StreamsGroupHeartbeatResponseData.TopicPartition>
standbyTopicPartitions = topicPartitions(standbyAndWarmupTasks,
configuredSubtopologies, metadataImage);
endpointToPartitions.setActivePartitions(activeTopicPartitions);
endpointToPartitions.setStandbyPartitions(standbyTopicPartitions);
return endpointToPartitions;
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManagerTest.java
index ffe05c3e40d..91b9594b81f 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/topics/EndpointToPartitionsManagerTest.java
@@ -113,6 +113,36 @@ class EndpointToPartitionsManagerTest {
assertTopicPartitionsAssigned(standbyPartitions, "Topic-B");
}
+ @Test
+ void testEndpointToPartitionsWithStandbyAndWarmupTasksInSameSubtopology() {
+ MetadataImage metadataImage = new MetadataImageBuilder()
+ .addTopic(Uuid.randomUuid(), "Topic-A", 3)
+ .addTopic(Uuid.randomUuid(), "Topic-B", 3)
+ .build();
+
+ when(streamsGroupMember.assignedTasks()).thenReturn(
+ new TasksTupleWithEpochs(
+ mkTasksPerSubtopologyWithCommonEpoch(0, mkEntry("0", Set.of(0,
1, 2))),
+ mkTasksPerSubtopology(mkEntry("1", Set.of(0))),
+ mkTasksPerSubtopology(mkEntry("1", Set.of(1, 2)))
+ )
+ );
+
when(streamsGroup.configuredTopology()).thenReturn(Optional.of(configuredTopology));
+ SortedMap<String, ConfiguredSubtopology> configuredSubtopologyMap =
new TreeMap<>();
+ configuredSubtopologyMap.put("0", configuredSubtopologyOne);
+ configuredSubtopologyMap.put("1", configuredSubtopologyTwo);
+
when(configuredTopology.subtopologies()).thenReturn(Optional.of(configuredSubtopologyMap));
+
+ StreamsGroupHeartbeatResponseData.EndpointToPartitions result =
+
EndpointToPartitionsManager.endpointToPartitions(streamsGroupMember,
responseEndpoint, streamsGroup, new
KRaftCoordinatorMetadataImage(metadataImage));
+
+ assertEquals(responseEndpoint, result.userEndpoint());
+ assertTopicPartitionsAssigned(result.activePartitions(), "Topic-A");
+ // The standby task and the warm-up tasks of subtopology 1 are merged
into a single standby entry.
+ assertEquals(1, result.standbyPartitions().size());
+ assertTopicPartitionsAssigned(result.standbyPartitions(), "Topic-B");
+ }
+
@Test
void
testEndpointToPartitionsExcludesPartitionsMissingFromSmallerSourceTopic() {
MetadataImage metadataImage = new MetadataImageBuilder()