Skip to content

Commit 7b9d847

Browse files
mjsaxMatthias J. Sax
authored andcommitted
KAFKA-20899: Fix "streams" StickyAssignor bug for standby tasks (#23096)
Broker side StickyTaskAssignor does not assign standby tasks correctly, because it does not cycle through all group member correctly to find the right candidate to host it. Reviewers: Mingi Cho (github:ChoMinGi), Bill Bejeck <bbejeck@apache.org>
1 parent 6b15c7e commit 7b9d847

2 files changed

Lines changed: 72 additions & 20 deletions

File tree

group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignor.java

Lines changed: 19 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import java.util.Iterator;
2828
import java.util.LinkedList;
2929
import java.util.Map;
30+
import java.util.Optional;
3031
import java.util.PriorityQueue;
3132
import java.util.Set;
3233
import java.util.stream.Collectors;
@@ -201,7 +202,7 @@ private void assignActive(final LinkedList<TaskId> activeTasks) {
201202
for (final Iterator<TaskId> it = activeTasks.iterator(); it.hasNext();) {
202203
final TaskId task = it.next();
203204
final ArrayList<Member> prevMembers = localState.standbyTaskToPrevMember.get(task);
204-
final Member prevMember = findPrevMemberWithLeastLoad(prevMembers, null);
205+
final Member prevMember = findPrevMemberWithLeastLoad(prevMembers, Optional.empty());
205206
if (prevMember != null) {
206207
final ProcessState processState = localState.processIdToState.get(prevMember.processId);
207208
if (hasUnfulfilledActiveTaskQuota(processState, prevMember)) {
@@ -279,33 +280,31 @@ private boolean assignStandbyToMemberWithLeastLoad(PriorityQueue<ProcessState> q
279280
*
280281
* @return Previous member with the least load that does not have the task, or null if no such member exists.
281282
*/
282-
private Member findPrevMemberWithLeastLoad(final ArrayList<Member> members, final TaskId taskId) {
283+
private Member findPrevMemberWithLeastLoad(final ArrayList<Member> members, final Optional<TaskId> standbyTaskId) {
283284
if (members == null || members.isEmpty()) {
284285
return null;
285286
}
286287

287-
Member candidate = members.get(0);
288-
final ProcessState candidateProcessState = localState.processIdToState.get(candidate.processId);
289-
double candidateProcessLoad = candidateProcessState.load();
290-
double candidateMemberLoad = candidateProcessState.memberToTaskCounts().get(candidate.memberId);
291-
for (int i = 1; i < members.size(); i++) {
292-
final Member member = members.get(i);
288+
Member candidate = null;
289+
double candidateProcessLoad = Double.MAX_VALUE;
290+
double candidateMemberLoad = Double.MAX_VALUE;
291+
for (final Member member : members) {
293292
final ProcessState processState = localState.processIdToState.get(member.processId);
293+
// A process that already owns a standby task (either as active or standby) cannot take it again
294+
if (standbyTaskId.isPresent() && processState.hasTask(standbyTaskId.get())) {
295+
continue;
296+
}
297+
294298
final double newProcessLoad = processState.load();
295-
if (newProcessLoad < candidateProcessLoad && (taskId == null || !processState.hasTask(taskId))) {
296-
final double newMemberLoad = processState.memberToTaskCounts().get(member.memberId);
297-
if (newMemberLoad < candidateMemberLoad) {
298-
candidateProcessLoad = newProcessLoad;
299-
candidateMemberLoad = newMemberLoad;
300-
candidate = member;
301-
}
299+
final double newMemberLoad = processState.memberToTaskCounts().get(member.memberId);
300+
if (candidate == null || (newProcessLoad < candidateProcessLoad && newMemberLoad < candidateMemberLoad)) {
301+
candidateProcessLoad = newProcessLoad;
302+
candidateMemberLoad = newMemberLoad;
303+
candidate = member;
302304
}
303305
}
304306

305-
if (taskId == null || !candidateProcessState.hasTask(taskId)) {
306-
return candidate;
307-
}
308-
return null;
307+
return candidate;
309308
}
310309

311310
private boolean hasUnfulfilledActiveTaskQuota(final ProcessState process, final Member member) {
@@ -339,7 +338,7 @@ private void assignStandby(final LinkedList<TaskId> standbyTasks) {
339338
// prev standby tasks
340339
final ArrayList<Member> prevStandbyMembers = localState.standbyTaskToPrevMember.get(task);
341340
if (prevStandbyMembers != null && !prevStandbyMembers.isEmpty()) {
342-
final Member prevStandbyMember = findPrevMemberWithLeastLoad(prevStandbyMembers, task);
341+
final Member prevStandbyMember = findPrevMemberWithLeastLoad(prevStandbyMembers, Optional.of(task));
343342
if (prevStandbyMember != null) {
344343
final ProcessState prevStandbyMemberProcessState = localState.processIdToState.get(prevStandbyMember.processId);
345344
if (hasUnfulfilledTaskQuota(prevStandbyMemberProcessState, prevStandbyMember)) {

group-coordinator/src/test/java/org/apache/kafka/coordinator/group/streams/assignor/StickyTaskAssignorTest.java

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import java.util.Arrays;
2525
import java.util.HashMap;
2626
import java.util.HashSet;
27+
import java.util.LinkedHashMap;
2728
import java.util.List;
2829
import java.util.Map;
2930
import java.util.NoSuchElementException;
@@ -1302,6 +1303,58 @@ public void shouldAssignStandbyTaskToPreviousOwnerBasedOnBelowQuotaCondition() {
13021303
}
13031304

13041305

1306+
@Test
1307+
public void shouldAssignStandbyToPreviousStandbyThatDoesNotHoldTheActiveTask() {
1308+
// starting assignment [active] [standby]:
1309+
// member1/process1: [] [2]
1310+
// member2/process2: [0] [2]
1311+
// member3/process3: [1] []
1312+
// member4/process4: [] []
1313+
//
1314+
// active 0,1 stay on their previous owner (stickiness)
1315+
// active 2 must go to member1/process1 as it has previous standby (standby->active promotion)
1316+
// -> it should not go to member2/process2 due to load balancing of active tasks
1317+
// standby 2 should stay on member2/process2 -> no reason to move it (load of member2 stays within capacity)
1318+
// standby 0,1: one *must* go to member4/process4 due to load balancing
1319+
// -> the other one can go anywhere but member2/process2 due to load balancing
1320+
//
1321+
// expected new assignment [active] [standby]
1322+
// ("expected" here means, based on the concrete implementation -- if this test fails with a different but
1323+
// still correct result -- as laid out above --, we should just update the test)
1324+
// member1/process1: [2] []
1325+
// member2/process2: [0] [2]
1326+
// member3/process3: [1] [0]
1327+
// member4/process4: [] [1]
1328+
final Map<String, AssignmentMemberSpec> members = new LinkedHashMap<>();
1329+
members.put("member1", createAssignmentMemberSpec("process1",
1330+
Map.of(), mkMap(mkEntry("test-subtopology", Set.of(2)))));
1331+
members.put("member2", createAssignmentMemberSpec("process2",
1332+
mkMap(mkEntry("test-subtopology", Set.of(0))), mkMap(mkEntry("test-subtopology", Set.of(2)))));
1333+
members.put("member3", createAssignmentMemberSpec("process3",
1334+
mkMap(mkEntry("test-subtopology", Set.of(1))), Map.of()));
1335+
members.put("member4", createAssignmentMemberSpec("process4"));
1336+
1337+
final GroupAssignment result = assignor.assign(
1338+
new GroupSpecImpl(members, mkMap(mkEntry(NUM_STANDBY_REPLICAS_CONFIG, "1"))),
1339+
new TopologyDescriberImpl(3, true, List.of("test-subtopology"))
1340+
);
1341+
1342+
assertEquals(Set.of(2), getActiveTasks(result, "test-subtopology", "member1"));
1343+
// if this is not empty, but only hold standby 0 or 1, still correct
1344+
assertEquals(List.of(), getAllStandbyTaskIds(result, "member1"));
1345+
1346+
assertEquals(Set.of(0), getActiveTasks(result, "test-subtopology", "member2"));
1347+
assertEquals(List.of(2), getAllStandbyTaskIds(result, "member2"));
1348+
1349+
assertEquals(Set.of(1), getActiveTasks(result, "test-subtopology", "member3"));
1350+
// if this is empty, or hold standby 1 instead, still correct
1351+
assertEquals(List.of(0), getAllStandbyTaskIds(result, "member3"));
1352+
1353+
assertEquals(Set.of(), getActiveTasks(result, "test-subtopology", "member4"));
1354+
// if this hold standby 0, or both standby 0 and 1, still correct
1355+
assertEquals(List.of(1), getAllStandbyTaskIds(result, "member4"));
1356+
}
1357+
13051358
private int getAllActiveTaskCount(GroupAssignment result, String... memberIds) {
13061359
int size = 0;
13071360
for (String memberId : memberIds) {

0 commit comments

Comments
 (0)