Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -553,7 +553,6 @@ private Map<String, ConnectorsAndTasks> performTaskRevocation(ConnectorsAndTasks
wl.worker(), wl.connectorsSize(), wl.tasksSize()));
}

Map<String, ConnectorsAndTasks> 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)) {
Expand All @@ -562,50 +561,68 @@ private Map<String, ConnectorsAndTasks> 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<String, ConnectorsAndTasks> 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<String> 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<String, ConnectorsAndTasks> revoking,
Collection<WorkerLoad> 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<ConnectorTaskId> 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<String, ExtendedAssignment> fillAssignments(Collection<String> members, short error,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

make sure there is no idle worker

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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -1237,6 +1238,66 @@ public void testDuplicatedAssignmentHandleWhenTheDuplicatedAssignmentsDeleted()
verify(coordinator, times(rebalanceNum)).lastCompletedGenerationId();
}

@Test
public void testComputeRevoked() {
Map<String, ConnectorsAndTasks> revoking = new HashMap<>();
Collection<WorkerLoad> 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();
}
Expand Down