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))

Reply via email to