diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java index f7fa55d8b3e9a..b68877870295e 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignor.java @@ -553,7 +553,6 @@ private Map performTaskRevocation(ConnectorsAndTasks wl.worker(), wl.connectorsSize(), wl.tasksSize())); } - Map revoking = new HashMap<>(); // If there are no new workers, or no existing workers to revoke tasks from return early // after logging the status if (!(newWorkersNum > 0 && existingWorkersNum > 0)) { @@ -562,50 +561,68 @@ private Map performTaskRevocation(ConnectorsAndTasks existingWorkersNum, newWorkersNum, totalWorkersNum); // This is intentionally empty but mutable, because the map is used to include deleted // connectors and tasks as well - return revoking; + return Collections.emptyMap(); } log.debug("Task revocation is required; workers with existing load: {} workers with " + "no load {} total workers {}", existingWorkersNum, newWorkersNum, totalWorkersNum); + Map revoking = new HashMap<>(); // We have at least one worker assignment (the leader itself) so totalWorkersNum can't be 0 log.debug("Previous rounded down (floor) average number of connectors per worker {}", totalActiveConnectorsNum / existingWorkersNum); - int floorConnectors = totalActiveConnectorsNum / totalWorkersNum; - int ceilConnectors = floorConnectors + ((totalActiveConnectorsNum % totalWorkersNum == 0) ? 0 : 1); - log.debug("New average number of connectors per worker rounded down (floor) {} and rounded up (ceil) {}", floorConnectors, ceilConnectors); - + computeRevoked( + revoking, + existingWorkers, + totalWorkersNum, + totalActiveConnectorsNum, + false + ); log.debug("Previous rounded down (floor) average number of tasks per worker {}", totalActiveTasksNum / existingWorkersNum); - int floorTasks = totalActiveTasksNum / totalWorkersNum; - int ceilTasks = floorTasks + ((totalActiveTasksNum % totalWorkersNum == 0) ? 0 : 1); - log.debug("New average number of tasks per worker rounded down (floor) {} and rounded up (ceil) {}", floorTasks, ceilTasks); - int numToRevoke; + computeRevoked( + revoking, + existingWorkers, + totalWorkersNum, + totalActiveTasksNum, + true + ); - for (WorkerLoad existing : existingWorkers) { - Iterator connectors = existing.connectors().iterator(); - numToRevoke = existing.connectorsSize() - ceilConnectors; - for (int i = existing.connectorsSize(); i > floorConnectors && numToRevoke > 0; --i, --numToRevoke) { - ConnectorsAndTasks resources = revoking.computeIfAbsent( - existing.worker(), - w -> new ConnectorsAndTasks.Builder().build()); - resources.connectors().add(connectors.next()); - } - } + return revoking; + } + static void computeRevoked(Map revoking, + Collection existingWorkers, + int numberOfWorkers, + int numberOfActiveTasks, + boolean forTask) { + int floor = numberOfActiveTasks / numberOfWorkers; + int numberOfCeilingMembers = numberOfActiveTasks % numberOfWorkers; + int ceil = numberOfCeilingMembers == 0 ? floor : floor + 1; for (WorkerLoad existing : existingWorkers) { - Iterator tasks = existing.tasks().iterator(); - numToRevoke = existing.tasksSize() - ceilTasks; - log.debug("Tasks on worker {} is higher than ceiling, so revoking {} tasks", existing, numToRevoke); - for (int i = existing.tasksSize(); i > floorTasks && numToRevoke > 0; --i, --numToRevoke) { + int currentSize = forTask ? existing.tasksSize() : existing.connectorsSize(); + int expectedSize; + // In order to reduce down-time, we have to avoid removing tasks from the nodes which are already in balance. + if (existingWorkers.size() == 1 || currentSize == 1) { + // this condition is used to deal with following specify cases + // 1) there is only one existent node so we can calculate the balanced numbers. + // 2) the node hosts only one task so it is unnecessary to remove the task from it. + expectedSize = ceil; + } else if (currentSize >= ceil && numberOfCeilingMembers > 0) { + expectedSize = ceil; + numberOfCeilingMembers--; + } else expectedSize = floor; + Iterator elements = forTask ? existing.tasks().iterator() : existing.connectors().iterator(); + int numToRevoke = currentSize - expectedSize; + while (elements.hasNext() && numToRevoke > 0) { ConnectorsAndTasks resources = revoking.computeIfAbsent( existing.worker(), w -> new ConnectorsAndTasks.Builder().build()); - resources.tasks().add(tasks.next()); + if (forTask) resources.tasks().add((ConnectorTaskId) elements.next()); + else resources.connectors().add((String) elements.next()); + numToRevoke--; } } - - return revoking; } private Map fillAssignments(Collection members, short error, diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/RebalanceSourceConnectorsIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/RebalanceSourceConnectorsIntegrationTest.java index b56a8faf35425..f2e60e4130aaa 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/RebalanceSourceConnectorsIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/RebalanceSourceConnectorsIntegrationTest.java @@ -362,7 +362,9 @@ private boolean assertConnectorAndTasksAreUniqueAndBalanced() { tasks.values().size(), tasks.values().stream().distinct().collect(Collectors.toList()).size()); assertTrue("Connectors are imbalanced: " + formatAssignment(connectors), maxConnectors - minConnectors < 2); + if (minConnectors > 1) assertEquals("Some workers have no connectors", connectors.size(), connect.workers().size()); assertTrue("Tasks are imbalanced: " + formatAssignment(tasks), maxTasks - minTasks < 2); + if (minTasks > 1) assertEquals("Some workers have no tasks", tasks.size(), connect.workers().size()); return true; } catch (Exception e) { log.error("Could not check connector state info.", e); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java index 0fe153132eb93..472dc68395d98 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/IncrementalCooperativeAssignorTest.java @@ -37,6 +37,7 @@ import java.util.AbstractMap.SimpleEntry; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -1237,6 +1238,66 @@ public void testDuplicatedAssignmentHandleWhenTheDuplicatedAssignmentsDeleted() verify(coordinator, times(rebalanceNum)).lastCompletedGenerationId(); } + @Test + public void testComputeRevoked() { + Map revoking = new HashMap<>(); + Collection existingWorkers = Arrays.asList( + new WorkerLoad.Builder("w0") + .with(Arrays.asList("c0", "c1"), + Arrays.asList(new ConnectorTaskId("c0", 0), + new ConnectorTaskId("c0", 1), + new ConnectorTaskId("c0", 2), + new ConnectorTaskId("c0", 3), + new ConnectorTaskId("c1", 0), + new ConnectorTaskId("c1", 1), + new ConnectorTaskId("c1", 2), + new ConnectorTaskId("c1", 3))) + .build(), + new WorkerLoad.Builder("w1") + .with(Arrays.asList("c2", "c3"), + Arrays.asList(new ConnectorTaskId("c2", 0), + new ConnectorTaskId("c2", 1), + new ConnectorTaskId("c2", 2), + new ConnectorTaskId("c2", 3), + new ConnectorTaskId("c3", 0), + new ConnectorTaskId("c3", 1), + new ConnectorTaskId("c3", 2), + new ConnectorTaskId("c3", 3))) + .build() + ); + + // test connectors + IncrementalCooperativeAssignor.computeRevoked( + revoking, + existingWorkers, + 3, + 4, + false + ); + // connectors distribution from (2, 2) -> (1, 2) + // total revoked: 1 + assertEquals(1, revoking.size()); + assertEquals(1, revoking.get("w1").connectors().size()); + + // test tasks + IncrementalCooperativeAssignor.computeRevoked( + revoking, + existingWorkers, + 3, + 16, + true + ); + + // tasks distribution from (8, 8) -> (6, 5) + // total revoked: 5 + assertEquals(2, revoking.size()); + int minNumberOfTasks = revoking.values().stream().mapToInt(c -> c.tasks().size()).min().orElse(0); + int maxNumberOfTasks = revoking.values().stream().mapToInt(c -> c.tasks().size()).max().orElse(0); + assertEquals(1, maxNumberOfTasks - minNumberOfTasks); + assertEquals(2, revoking.get("w0").tasks().size()); + assertEquals(3, revoking.get("w1").tasks().size()); + } + private WorkerLoad emptyWorkerLoad(String worker) { return new WorkerLoad.Builder(worker).build(); }