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 68903c7a7d7 KAFKA-20665: Add refiner scaffolding for "streams" warmup
tasks (#22658)
68903c7a7d7 is described below
commit 68903c7a7d72782585729cec17fdbf24809e360d
Author: Matthias J. Sax <[email protected]>
AuthorDate: Thu Jul 2 17:01:07 2026 -0700
KAFKA-20665: Add refiner scaffolding for "streams" warmup tasks (#22658)
Part of KIP-1071.
After an assignment is computed, we need to refine it to apply warmup
tasks. This PR adds the basic control flow into the refiner method; The
actual refiner logic will be added in a follow up PR; for now, we just
return the unmodified assignment.
The refiner needs access to the whole group assignment, not just the
assignment of a single member. Thus, the PR contains some minor
refactoring for this.
Reviewers: Sean Quah <[email protected]>
---
.../coordinator/group/GroupMetadataManager.java | 125 +++++++++++++--------
...sGroupStaticMemberGroupMetadataManagerTest.java | 6 +
2 files changed, 87 insertions(+), 44 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 97a49040f56..9643099e088 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
@@ -152,6 +152,7 @@ import
org.apache.kafka.coordinator.group.modern.share.ShareGroup.InitMapValue;
import
org.apache.kafka.coordinator.group.modern.share.ShareGroup.ShareGroupStatePartitionMetadataInfo;
import
org.apache.kafka.coordinator.group.modern.share.ShareGroupAssignmentBuilder;
import org.apache.kafka.coordinator.group.modern.share.ShareGroupMember;
+import org.apache.kafka.coordinator.group.streams.MemberTaskOffsets;
import
org.apache.kafka.coordinator.group.streams.StreamsCoordinatorRecordHelpers;
import org.apache.kafka.coordinator.group.streams.StreamsGroup;
import org.apache.kafka.coordinator.group.streams.StreamsGroupDescribeResult;
@@ -309,23 +310,6 @@ public class GroupMetadataManager {
group.targetAssignment(member.memberId())
);
}
-
- private static UpdateTargetAssignmentResult<TasksTuple>
fromLastTargetAssignment(
- StreamsGroup group,
- Optional<StreamsGroupMember> member
- ) {
- if (member.isPresent()) {
- return new UpdateTargetAssignmentResult<>(
- group.assignmentEpoch(),
- group.targetAssignment(member.get().memberId(),
member.get().instanceId())
- );
- } else {
- return new UpdateTargetAssignmentResult<>(
- group.assignmentEpoch(),
- TasksTuple.EMPTY
- );
- }
- }
}
public static class Builder {
@@ -2259,7 +2243,7 @@ public class GroupMetadataManager {
// 4. Update the target assignment if the group epoch is larger than
the target assignment epoch or a static member
// replaces an existing static member.
// The delta between the existing and the new target assignment is
persisted to the partition.
- UpdateTargetAssignmentResult<TasksTuple> updateTargetAssignmentResult
= maybeUpdateStreamsTargetAssignment(
+ UpdateTargetAssignmentResult<Map<String, TasksTuple>>
updateTargetAssignmentResult = maybeUpdateStreamsTargetAssignment(
group,
groupEpoch,
Optional.of(updatedMember),
@@ -2270,7 +2254,17 @@ public class GroupMetadataManager {
currentAssignmentConfigs
);
- // 5. Reconcile the member's assignment with the target assignment if
the member is not
+ // 4b. Refine the target assignment into the intermediate assignment
(with warm-up tasks) the member should be
+ // reconciled toward. Runs on every heartbeat, before reconciliation.
No-op for now.
+ TasksTuple refinedTarget = refine(
+ updatedMember,
+ updateTargetAssignmentResult.targetAssignment,
+ group.taskOffsets(),
+ streamsGroupNumWarmupReplicas(group.groupId()),
+ streamsGroupAcceptableRecoveryLag(group.groupId())
+ );
+
+ // 5. Reconcile the member's assignment with the (refined) target
assignment if the member is not
// fully reconciled yet.
updatedMember = maybeReconcile(
groupId,
@@ -2279,7 +2273,7 @@ public class GroupMetadataManager {
group::currentStandbyTaskProcessIds,
group::currentWarmupTaskProcessIds,
updateTargetAssignmentResult.targetAssignmentEpoch(),
- updateTargetAssignmentResult.targetAssignment(),
+ refinedTarget,
ownedActiveTasks,
ownedStandbyTasks,
ownedWarmupTasks,
@@ -4322,9 +4316,9 @@ public class GroupMetadataManager {
* @param metadataImage The metadata image.
* @param records The list to accumulate any new records.
* @param returnedStatus A mutable collection of status to be
returned in the response.
- * @return The new target assignment for the updated member, or EMPTY if
no member specified.
+ * @return The target assignment epoch and the full per-member target
assignment.
*/
- private UpdateTargetAssignmentResult<TasksTuple>
maybeUpdateStreamsTargetAssignment(
+ private UpdateTargetAssignmentResult<Map<String, TasksTuple>>
maybeUpdateStreamsTargetAssignment(
StreamsGroup group,
int groupEpoch,
Optional<StreamsGroupMember> updatedMember,
@@ -4342,15 +4336,26 @@ public class GroupMetadataManager {
.setStatusDetail("Assignment delayed due to the configured
initial rebalance delay.")
));
- return new UpdateTargetAssignmentResult<>(
- group.assignmentEpoch(),
- TasksTuple.EMPTY
- );
+ return new UpdateTargetAssignmentResult<>(group.assignmentEpoch(),
Map.of());
}
+ // Apply the unwritten membership change (e.g. the relabelling of a
replacing static member) to the
+ // target assignment, so both the assignor and the refiner see a
consistent view. The relabelling
+ // record is only queued during this heartbeat and not yet replayed
into the in-memory assignment.
+ UpdatedMembersAndTargetAssignmentView<StreamsGroupMember, TasksTuple>
updatedMembersAndTargetAssignment =
+ new UpdatedMembersAndTargetAssignmentView<>(
+ group.members(),
+ group.staticMembers(),
+ group.targetAssignment(),
+ m -> m.instanceId().orElse(null)
+ );
+ updatedMember.ifPresent(member ->
+
updatedMembersAndTargetAssignment.addOrUpdateMember(member.memberId(), member)
+ );
+
if (group.assignmentEpoch() >= groupEpoch) {
// The assignment is up to date.
- return
UpdateTargetAssignmentResult.fromLastTargetAssignment(group, updatedMember);
+ return new UpdateTargetAssignmentResult<>(group.assignmentEpoch(),
updatedMembersAndTargetAssignment.targetAssignment());
}
boolean canComputeNextTargetAssignment =
canComputeNextTargetAssignment(
@@ -4365,22 +4370,11 @@ public class GroupMetadataManager {
.setStatusDetail("Assignment delayed due to the configured
assignment interval.")
));
- return
UpdateTargetAssignmentResult.fromLastTargetAssignment(group, updatedMember);
+ return new UpdateTargetAssignmentResult<>(group.assignmentEpoch(),
updatedMembersAndTargetAssignment.targetAssignment());
}
TaskAssignor assignor = streamsGroupAssignor(group.groupId());
try {
- UpdatedMembersAndTargetAssignmentView<StreamsGroupMember,
TasksTuple> updatedMembersAndTargetAssignment =
- new UpdatedMembersAndTargetAssignmentView<>(
- group.members(),
- group.staticMembers(),
- group.targetAssignment(),
- m -> m.instanceId().orElse(null)
- );
- updatedMember.ifPresent(member ->
-
updatedMembersAndTargetAssignment.addOrUpdateMember(member.memberId(), member)
- );
-
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder
assignmentResultBuilder =
new
org.apache.kafka.coordinator.group.streams.TargetAssignmentBuilder(
group.groupId(),
@@ -4410,11 +4404,7 @@ public class GroupMetadataManager {
records.addAll(assignmentResult.records());
- return new UpdateTargetAssignmentResult<>(
- groupEpoch,
- updatedMember.map(member ->
assignmentResult.targetAssignment().get(member.memberId()))
- .orElse(TasksTuple.EMPTY)
- );
+ return new UpdateTargetAssignmentResult<>(groupEpoch,
assignmentResult.targetAssignment());
} catch (TaskAssignorException ex) {
String msg = String.format("Failed to compute a new target
assignment for epoch %d: %s",
groupEpoch, ex.getMessage());
@@ -4423,6 +4413,44 @@ public class GroupMetadataManager {
}
}
+
+ /**
+ * Refines the task assignor's target assignment into the
<em>intermediate</em> assignment that
+ * the reconciler ({@link
org.apache.kafka.coordinator.group.streams.CurrentAssignmentBuilder}) converges
members toward.
+ * <p>
+ * The intermediate assignment is the current assignment with warm-up
tasks inserted (for later promotion to active,
+ * based on per-member changelog lag), so that a task is moved to a new
owner only once that owner has caught up.
+ * It is held in memory only and is never persisted: the assignor's target
assignment remains the persisted source of
+ * truth (and the intermediate is reconstructed from the persisted
current-assignment records after a coordinator failover).
+ * <p>
+ * The refiner is invoked on every heartbeat, before reconciliation, so it
can react to all the inputs that can change
+ * the intermediate assignment between reassignments — newly reported task
offsets (a warm-up may have become hot),
+ * {@code num.warmup.replicas} / {@code acceptable.recovery.lag} config
changes, and members acknowledging task
+ * revocation/restoration (advancing an in-flight migration).
+ *
+ * @param member
+ * The member to produce the refined (intermediate) assignment for.
+ * @param targetAssignment
+ * All members' target assignments (group context).
+ * @param taskOffsets
+ * The latest per-member changelog offsets/end-offsets reported via
heartbeats (group context).
+ * @param numWarmupReplicas
+ * The configured maximum number of warm-up replicas.
+ * @param acceptableRecoveryLag
+ * The lag at or below which a warm-up is considered caught up and
can be promoted.
+ *
+ * @return The member's intermediate assignment tuple.
+ */
+ private static TasksTuple refine(
+ final StreamsGroupMember member,
+ final Map<String, TasksTuple> targetAssignment,
+ final Map<String, MemberTaskOffsets> taskOffsets,
+ final int numWarmupReplicas,
+ final long acceptableRecoveryLag
+ ) {
+ return targetAssignment.getOrDefault(member.memberId(),
TasksTuple.EMPTY);
+ }
+
/**
* Fires the initial rebalance for a streams group when the delay timer
expires.
* Computes and persists target assignment for all members if conditions
are met.
@@ -9528,6 +9556,15 @@ public class GroupMetadataManager {
.orElse(config.streamsGroupAcceptableRecoveryLag());
}
+ /**
+ * Get the number of warmup replicas of the provided streams group.
+ */
+ private int streamsGroupNumWarmupReplicas(String groupId) {
+ Optional<GroupConfig> groupConfig =
groupConfigManager.groupConfig(groupId);
+ return groupConfig.flatMap(GroupConfig::streamsNumWarmupReplicas)
+ .orElse(config.streamsGroupNumWarmupReplicas());
+ }
+
/**
* Get the initial rebalance delay of the provided streams group.
*/
diff --git
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
index 3584313d7ca..48fe6bb834a 100644
---
a/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
+++
b/group-coordinator/src/test/java/org/apache/kafka/coordinator/group/StreamsGroupStaticMemberGroupMetadataManagerTest.java
@@ -611,6 +611,12 @@ class StreamsGroupStaticMemberGroupMetadataManagerTest {
assertEquals(rejoinMemberId, result.response().data().memberId());
assertEquals(groupEpoch, result.response().data().memberEpoch());
+ // The static member rejoined with a new member id while the
assignment is up to date, so the
+ // target assignment is reused from the persisted state. Its
relabelling to the new member id
+ // is queued but not yet replayed, so the rejoining member must still
get its assignment back
+ // (not an empty one).
+ assertEquals(topic.responseTasks(0, 1, 2, 3),
result.response().data().activeTasks());
+
Optional<StreamsGroupMemberMetadataValue> updatedMemberMetadataValue =
result.records().stream()
.filter(record -> record.key() instanceof
StreamsGroupMemberMetadataKey)
.filter(record -> ((StreamsGroupMemberMetadataKey)
record.key()).memberId().equals(rejoinMemberId))