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 774e0ab6b1b MINOR: "streams" TaskAssignors must be thread-safe (#22930)
774e0ab6b1b is described below

commit 774e0ab6b1b93cfa8806dd38ca56690574cde783
Author: Matthias J. Sax <[email protected]>
AuthorDate: Mon Jul 27 07:38:54 2026 -0700

    MINOR: "streams" TaskAssignors must be thread-safe (#22930)
    
    With KIP-1357, task assignors are instantiated differently and now must
    be thread safe.
    
    This PR updates the interface JavaDocs accordingly.
    
    Additionally, it fixes StickyTaskAssignor, which currently keeps local
    state in an instance field, what is not thread safe. This PR refactors
    StickyTaskAssignor accordingly.
    
    Reviewers: TengYao Chi <[email protected]>
---
 .../group/api/streams/assignor/TaskAssignor.java   |   2 +
 .../group/streams/assignor/StickyTaskAssignor.java | 104 ++++++++++++---------
 2 files changed, 64 insertions(+), 42 deletions(-)

diff --git 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java
 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java
index f6532b8c840..def7dc62356 100644
--- 
a/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java
+++ 
b/group-coordinator/group-coordinator-api/src/main/java/org/apache/kafka/coordinator/group/api/streams/assignor/TaskAssignor.java
@@ -21,6 +21,8 @@ import org.apache.kafka.common.annotation.InterfaceStability;
 
 /**
  * Server side task assignor used by streams groups.
+ *
+ * <p>Implementations must be thread-safe.
  */
 @InterfaceAudience.Public
 @InterfaceStability.Evolving
diff --git 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java
 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java
index b1686b426aa..529b07fd4fd 100644
--- 
a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java
+++ 
b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java
@@ -45,9 +45,6 @@ public class StickyTaskAssignor implements TaskAssignor {
     private static final String STICKY_ASSIGNOR_NAME = "sticky";
     private static final Logger log = 
LoggerFactory.getLogger(StickyTaskAssignor.class);
 
-    private LocalState localState;
-
-
     @Override
     public String name() {
         return STICKY_ASSIGNOR_NAME;
@@ -60,25 +57,30 @@ public class StickyTaskAssignor implements TaskAssignor {
 
     @Override
     public GroupAssignment assign(final GroupSpec groupSpec, final 
TopologyDescriber topologyDescriber) throws TaskAssignorException {
-        initialize(groupSpec, topologyDescriber);
-        final GroupAssignment assignments =  doAssign(groupSpec, 
topologyDescriber);
-        localState = null;
-        return assignments;
+        return doAssign(
+            initialize(groupSpec, topologyDescriber),
+            groupSpec,
+            topologyDescriber
+        );
     }
 
-    private GroupAssignment doAssign(final GroupSpec groupSpec, final 
TopologyDescriber topologyDescriber) {
+    private static GroupAssignment doAssign(
+        final LocalState localState,
+        final GroupSpec groupSpec,
+        final TopologyDescriber topologyDescriber
+    ) {
         final LinkedList<TaskId> activeTasks = taskIds(topologyDescriber, 
true);
-        assignActive(activeTasks);
+        assignActive(localState, activeTasks);
 
         if (localState.numStandbyReplicas > 0) {
             final LinkedList<TaskId> statefulTasks = 
taskIds(topologyDescriber, false);
-            assignStandby(statefulTasks);
+            assignStandby(localState, statefulTasks);
         }
 
-        return buildGroupAssignment(groupSpec.memberIds());
+        return buildGroupAssignment(localState, groupSpec.memberIds());
     }
 
-    private LinkedList<TaskId> taskIds(final TopologyDescriber 
topologyDescriber, final boolean isActive) {
+    private static LinkedList<TaskId> taskIds(final TopologyDescriber 
topologyDescriber, final boolean isActive) {
         final LinkedList<TaskId> ret = new LinkedList<>();
         for (final String subtopology : topologyDescriber.subtopologies()) {
             if (isActive || topologyDescriber.isStateful(subtopology)) {
@@ -91,8 +93,8 @@ public class StickyTaskAssignor implements TaskAssignor {
         return ret;
     }
 
-    private void initialize(final GroupSpec groupSpec, final TopologyDescriber 
topologyDescriber) {
-        localState = new LocalState();
+    private static LocalState initialize(final GroupSpec groupSpec, final 
TopologyDescriber topologyDescriber) {
+        final LocalState localState = new LocalState();
         localState.numStandbyReplicas =
             groupSpec.configs().isEmpty() ? 0
                 : 
Integer.parseInt(groupSpec.configs().get("num.standby.replicas"));
@@ -141,9 +143,10 @@ public class StickyTaskAssignor implements TaskAssignor {
                 }
             }
         }
+        return localState;
     }
 
-    private GroupAssignment buildGroupAssignment(final Collection<String> 
members) {
+    private static GroupAssignment buildGroupAssignment(final LocalState 
localState, final Collection<String> members) {
         final Map<String, MemberAssignment> memberAssignments = new 
HashMap<>();
 
         final Map<String, Set<TaskId>> activeTasksAssignments = 
localState.processIdToState.entrySet().stream()
@@ -176,7 +179,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         return new GroupAssignment(memberAssignments);
     }
 
-    private Map<String, Set<Integer>> toCompactedTaskIds(final Set<TaskId> 
taskIds) {
+    private static Map<String, Set<Integer>> toCompactedTaskIds(final 
Set<TaskId> taskIds) {
         final Map<String, Set<Integer>> ret = new HashMap<>();
         for (final TaskId taskId : taskIds) {
             ret.putIfAbsent(taskId.subtopologyId(), new HashSet<>());
@@ -185,7 +188,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         return ret;
     }
 
-    private void assignActive(final LinkedList<TaskId> activeTasks) {
+    private static void assignActive(final LocalState localState, final 
LinkedList<TaskId> activeTasks) {
 
         // Assuming our current assignment pairs same partitions 
(range-based), we want to sort by partition first
         
activeTasks.sort(Comparator.comparing(TaskId::partition).thenComparing(TaskId::subtopologyId));
@@ -196,10 +199,10 @@ public class StickyTaskAssignor implements TaskAssignor {
             final Member prevMember = 
localState.activeTaskToPrevMember.get(task);
             if (prevMember != null) {
                 final ProcessState processState = 
localState.processIdToState.get(prevMember.processId);
-                if (hasUnfulfilledActiveTaskQuota(processState, prevMember)) {
+                if (hasUnfulfilledActiveTaskQuota(localState, processState, 
prevMember)) {
                     int newActiveTasks = 
processState.addTask(prevMember.memberId, task, true);
-                    maybeUpdateActiveTasksPerMember(newActiveTasks);
-                    maybeUpdateTotalTasksPerMember(newActiveTasks);
+                    maybeUpdateActiveTasksPerMember(localState, 
newActiveTasks);
+                    maybeUpdateTotalTasksPerMember(localState, newActiveTasks);
                     it.remove();
                 }
             }
@@ -209,13 +212,13 @@ public class StickyTaskAssignor implements TaskAssignor {
         for (final Iterator<TaskId> it = activeTasks.iterator(); 
it.hasNext();) {
             final TaskId task = it.next();
             final ArrayList<Member> prevMembers = 
localState.standbyTaskToPrevMember.get(task);
-            final Member prevMember = findPrevMemberWithLeastLoad(prevMembers, 
null);
+            final Member prevMember = findPrevMemberWithLeastLoad(localState, 
prevMembers, null);
             if (prevMember != null) {
                 final ProcessState processState = 
localState.processIdToState.get(prevMember.processId);
-                if (hasUnfulfilledActiveTaskQuota(processState, prevMember)) {
+                if (hasUnfulfilledActiveTaskQuota(localState, processState, 
prevMember)) {
                     int newActiveTasks = 
processState.addTask(prevMember.memberId, task, true);
-                    maybeUpdateActiveTasksPerMember(newActiveTasks);
-                    maybeUpdateTotalTasksPerMember(newActiveTasks);
+                    maybeUpdateActiveTasksPerMember(localState, 
newActiveTasks);
+                    maybeUpdateTotalTasksPerMember(localState, newActiveTasks);
                     it.remove();
                 }
             }
@@ -234,8 +237,8 @@ public class StickyTaskAssignor implements TaskAssignor {
             }
             final int newTaskCount = 
processWithLeastLoad.addTaskToLeastLoadedMember(task, true);
             if (newTaskCount != -1) {
-                maybeUpdateActiveTasksPerMember(newTaskCount);
-                maybeUpdateTotalTasksPerMember(newTaskCount);
+                maybeUpdateActiveTasksPerMember(localState, newTaskCount);
+                maybeUpdateTotalTasksPerMember(localState, newTaskCount);
             } else {
                 throw new TaskAssignorException(String.format("No member 
available to assign active task %s.", task));
             }
@@ -243,7 +246,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         }
     }
 
-    private void maybeUpdateActiveTasksPerMember(final int activeTasksNo) {
+    private static void maybeUpdateActiveTasksPerMember(final LocalState 
localState, final int activeTasksNo) {
         if (activeTasksNo == localState.activeTasksPerMember) {
             localState.totalMembersWithActiveTaskCapacity--;
             localState.totalActiveTasks -= activeTasksNo;
@@ -251,7 +254,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         }
     }
 
-    private void maybeUpdateTotalTasksPerMember(final int taskNo) {
+    private static void maybeUpdateTotalTasksPerMember(final LocalState 
localState, final int taskNo) {
         if (taskNo == localState.totalTasksPerMember) {
             localState.totalMembersWithTaskCapacity--;
             localState.totalTasks -= taskNo;
@@ -259,7 +262,11 @@ public class StickyTaskAssignor implements TaskAssignor {
         }
     }
 
-    private boolean 
assignStandbyToMemberWithLeastLoad(PriorityQueue<ProcessState> queue, TaskId 
taskId) {
+    private static boolean assignStandbyToMemberWithLeastLoad(
+        final LocalState localState,
+        final PriorityQueue<ProcessState> queue,
+        final TaskId taskId
+    ) {
         final ProcessState processWithLeastLoad = queue.poll();
         if (processWithLeastLoad == null) {
             return false;
@@ -269,10 +276,10 @@ public class StickyTaskAssignor implements TaskAssignor {
             final int newTaskCount = 
processWithLeastLoad.addTaskToLeastLoadedMember(taskId, false);
             if (newTaskCount != -1) {
                 found = true;
-                maybeUpdateTotalTasksPerMember(newTaskCount);
+                maybeUpdateTotalTasksPerMember(localState, newTaskCount);
             }
         } else if (!queue.isEmpty()) {
-            found = assignStandbyToMemberWithLeastLoad(queue, taskId);
+            found = assignStandbyToMemberWithLeastLoad(localState, queue, 
taskId);
         }
         queue.add(processWithLeastLoad); // Add it back to the queue after 
updating its state
         return found;
@@ -281,13 +288,18 @@ public class StickyTaskAssignor implements TaskAssignor {
     /**
      * Finds the previous member with the least load for a given task.
      *
+     * @param localState The state of the assignment in progress.
      * @param members The list of previous members owning the task.
      * @param taskId  The taskId, to check if the previous member already has 
the task. Can be null, if we assign it
      *                for the first time (e.g., during active task assignment).
      *
      * @return Previous member with the least load that does not have the 
task, or null if no such member exists.
      */
-    private Member findPrevMemberWithLeastLoad(final ArrayList<Member> 
members, final TaskId taskId) {
+    private static Member findPrevMemberWithLeastLoad(
+        final LocalState localState,
+        final ArrayList<Member> members,
+        final TaskId taskId
+    ) {
         if (members == null || members.isEmpty()) {
             return null;
         }
@@ -316,15 +328,23 @@ public class StickyTaskAssignor implements TaskAssignor {
         return null;
     }
 
-    private boolean hasUnfulfilledActiveTaskQuota(final ProcessState process, 
final Member member) {
+    private static boolean hasUnfulfilledActiveTaskQuota(
+        final LocalState localState,
+        final ProcessState process,
+        final Member member
+    ) {
         return process.memberToTaskCounts().get(member.memberId) < 
localState.activeTasksPerMember;
     }
 
-    private boolean hasUnfulfilledTaskQuota(final ProcessState process, final 
Member member) {
+    private static boolean hasUnfulfilledTaskQuota(
+        final LocalState localState,
+        final ProcessState process,
+        final Member member
+    ) {
         return process.memberToTaskCounts().get(member.memberId) < 
localState.totalTasksPerMember;
     }
 
-    private void assignStandby(final LinkedList<TaskId> standbyTasks) {
+    private static void assignStandby(final LocalState localState, final 
LinkedList<TaskId> standbyTasks) {
         final ArrayList<StandbyToAssign> toLeastLoaded = new 
ArrayList<>(standbyTasks.size() * localState.numStandbyReplicas);
         
         // Assuming our current assignment is range-based, we want to sort by 
partition first.
@@ -337,9 +357,9 @@ public class StickyTaskAssignor implements TaskAssignor {
                 final Member prevActiveMember = 
localState.activeTaskToPrevMember.get(task);
                 if (prevActiveMember != null) {
                     final ProcessState prevActiveMemberProcessState = 
localState.processIdToState.get(prevActiveMember.processId);
-                    if (!prevActiveMemberProcessState.hasTask(task) && 
hasUnfulfilledTaskQuota(prevActiveMemberProcessState, prevActiveMember)) {
+                    if (!prevActiveMemberProcessState.hasTask(task) && 
hasUnfulfilledTaskQuota(localState, prevActiveMemberProcessState, 
prevActiveMember)) {
                         int newTaskCount = 
prevActiveMemberProcessState.addTask(prevActiveMember.memberId, task, false);
-                        maybeUpdateTotalTasksPerMember(newTaskCount);
+                        maybeUpdateTotalTasksPerMember(localState, 
newTaskCount);
                         continue;
                     }
                 }
@@ -347,12 +367,12 @@ public class StickyTaskAssignor implements TaskAssignor {
                 // prev standby tasks
                 final ArrayList<Member> prevStandbyMembers = 
localState.standbyTaskToPrevMember.get(task);
                 if (prevStandbyMembers != null && 
!prevStandbyMembers.isEmpty()) {
-                    final Member prevStandbyMember = 
findPrevMemberWithLeastLoad(prevStandbyMembers, task);
+                    final Member prevStandbyMember = 
findPrevMemberWithLeastLoad(localState, prevStandbyMembers, task);
                     if (prevStandbyMember != null) {
                         final ProcessState prevStandbyMemberProcessState = 
localState.processIdToState.get(prevStandbyMember.processId);
-                        if 
(hasUnfulfilledTaskQuota(prevStandbyMemberProcessState, prevStandbyMember)) {
+                        if (hasUnfulfilledTaskQuota(localState, 
prevStandbyMemberProcessState, prevStandbyMember)) {
                             int newTaskCount = 
prevStandbyMemberProcessState.addTask(prevStandbyMember.memberId, task, false);
-                            maybeUpdateTotalTasksPerMember(newTaskCount);
+                            maybeUpdateTotalTasksPerMember(localState, 
newTaskCount);
                             continue;
                         }
                     }
@@ -371,7 +391,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         processByLoad.addAll(localState.processIdToState.values());
         for (final StandbyToAssign toAssign : toLeastLoaded) {
             for (int i = 0; i < toAssign.remainingReplicas; i++) {
-                if (!assignStandbyToMemberWithLeastLoad(processByLoad, 
toAssign.taskId)) {
+                if (!assignStandbyToMemberWithLeastLoad(localState, 
processByLoad, toAssign.taskId)) {
                     log.warn("{} There is not enough available capacity. " +
                             "You should increase the number of threads and/or 
application instances to maintain the requested number of standby replicas.",
                         errorMessage(localState.numStandbyReplicas, i, 
toAssign.taskId));
@@ -381,7 +401,7 @@ public class StickyTaskAssignor implements TaskAssignor {
         }
     }
 
-    private String errorMessage(final int numStandbyReplicas, final int i, 
final TaskId task) {
+    private static String errorMessage(final int numStandbyReplicas, final int 
i, final TaskId task) {
         return "Unable to assign " + (numStandbyReplicas - i) +
             " of " + numStandbyReplicas + " standby tasks for task [" + task + 
"].";
     }

Reply via email to