From 007b12909b62491ac5ea1a7cba9af58f3d1d6a33 Mon Sep 17 00:00:00 2001 From: Bob Haddleton Date: Mon, 8 Feb 2021 12:00:54 -0600 Subject: [PATCH 1/2] improve partition assignment balance (#93) --- faust/assignor/copartitioned_assignor.py | 47 ++++++++++++++----- .../assignor/test_copartitioned_assignor.py | 15 +++++- 2 files changed, 47 insertions(+), 15 deletions(-) diff --git a/faust/assignor/copartitioned_assignor.py b/faust/assignor/copartitioned_assignor.py index 58785bfcd..1c762ce85 100644 --- a/faust/assignor/copartitioned_assignor.py +++ b/faust/assignor/copartitioned_assignor.py @@ -1,6 +1,6 @@ """Copartitioned Assignor.""" from itertools import cycle -from math import ceil +from math import ceil, floor from typing import Iterable, Iterator, MutableMapping, Optional, Sequence, Set from mode.utils.typing import Counter @@ -17,12 +17,12 @@ class CopartitionedAssignor: The assignment is sticky which uses the following heuristics: - - Maintain existing assignments as long as within capacity for each client - - Assign actives to standbys when possible (within capacity) - - Assign in order to fill capacity of the clients + - Maintain existing assignments as long as within max_capacity for each client + - Assign actives to standbys when possible (within max_capacity) + - Assign in order to fill max_capacity of the clients We optimize for not over utilizing resources instead of under-utilizing - resources. This results in a balanced assignment when capacity is the + resources. This results in a balanced assignment when max_capacity is the default value which is ``ceil(num partitions / num clients)`` Notes: @@ -30,7 +30,8 @@ class CopartitionedAssignor: for the desired `replication`. """ - capacity: int + max_capacity: int + min_capacity: int num_partitions: int replicas: int topics: Set[str] @@ -50,22 +51,31 @@ def __init__( assert self._num_clients, "Should assign to at least 1 client" self.num_partitions = num_partitions self.replicas = min(replicas, self._num_clients - 1) - self.capacity = ( + self.max_capacity = ( int(ceil(float(self.num_partitions) / self._num_clients)) if capacity is None else capacity ) + self.min_capacity = ( + int(floor(float(self.num_partitions) / self._num_clients)) + if capacity is None + else capacity + ) self.topics = set(topics) assert ( - self.capacity * self._num_clients >= self.num_partitions - ), "Not enough capacity" + self.max_capacity * self._num_clients >= self.num_partitions + ), "Not enough max_capacity" self._client_assignments = cluster_asgn def get_assignment(self) -> MutableMapping[str, CopartitionedAssignment]: + # If there are no replicas required, try to ensure a balanced distribution with + # at least one partition per client by removing partitions down to the + # min_capacity for each client + capacity = self.min_capacity if self.replicas == 0 else self.max_capacity for copartitioned in self._client_assignments.values(): - copartitioned.unassign_extras(self.capacity, self.replicas) + copartitioned.unassign_extras(capacity, self.replicas) self._assign(active=True) self._assign(active=False) return self._client_assignments @@ -92,7 +102,7 @@ def _assigned_partition_counts(self, active: bool) -> Counter[int]: ) def _get_client_limit(self, active: bool) -> int: - return self.capacity * self._total_assigns_per_partition(active) + return self.max_capacity * self._total_assigns_per_partition(active) def _total_assigns_per_partition(self, active: bool) -> int: return 1 if active else self.replicas @@ -172,17 +182,28 @@ def _find_round_robin_assignable( def _assign_round_robin(self, unassigned: Iterable[int], active: bool) -> None: # We do round robin assignment as follows: + # - Sort the candidate clients by the number of partitions already assigned, + # to improve the overall balance of the assignment process # - For actives, we first try to assign to a standby # - For standby, we offset the start for round robin to evenly # distribute standbys for colocated actives - # - We do round robin + # - We do round robin over the sorted clients # - If no assignment found, it must be a standby and the only unfilled # client(s) must be actives/standbys for the partition # - If no assignment found, we unassign and arbitrary partition from a # filled assignment such that the partition can be assigned to it # - This guarantees eventual assignment of all partitions client_limit = self._get_client_limit(active) - candidates = cycle(self._client_assignments.values()) + if active: + candidates_list = sorted( + list(self._client_assignments.values()), key=lambda x: len(x.actives) + ) + else: + candidates_list = sorted( + list(self._client_assignments.values()), key=lambda x: len(x.standbys) + ) + + candidates = cycle(iter(candidates_list)) unassigned = list(unassigned) while unassigned: partition = unassigned.pop(0) diff --git a/tests/meticulous/assignor/test_copartitioned_assignor.py b/tests/meticulous/assignor/test_copartitioned_assignor.py index 3156e5e3a..8179012e1 100644 --- a/tests/meticulous/assignor/test_copartitioned_assignor.py +++ b/tests/meticulous/assignor/test_copartitioned_assignor.py @@ -41,6 +41,11 @@ def is_valid( return True +def clients_balanced(new: MutableMapping[str, CopartitionedAssignment]) -> bool: + assert all(new[client].actives for client in new) + return True + + def client_addition_sticky( old: MutableMapping[str, CopartitionedAssignment], new: MutableMapping[str, CopartitionedAssignment], @@ -120,7 +125,10 @@ def test_add_new_clients(partitions, replicas, num_clients, num_additional_clien ).get_assignment() assert is_valid(new_assignments, partitions, replicas) - assert client_addition_sticky(old_assignments, new_assignments) + if replicas > 0: + assert client_addition_sticky(old_assignments, new_assignments) + elif partitions and partitions >= num_clients + num_additional_clients: + assert clients_balanced(new_assignments) @given( @@ -156,4 +164,7 @@ def test_remove_clients(partitions, replicas, num_clients, num_removal_clients): ).get_assignment() assert is_valid(new_assignments, partitions, replicas) - assert client_removal_sticky(old_assignments, new_assignments) + if replicas > 0: + assert client_removal_sticky(old_assignments, new_assignments) + elif partitions >= num_clients - num_removal_clients: + assert clients_balanced(new_assignments) From d05caccc7dce430955b6094ff207eb4227cd18ed Mon Sep 17 00:00:00 2001 From: Bob Haddleton Date: Tue, 9 Feb 2021 11:30:07 -0600 Subject: [PATCH 2/2] restrict partition balancing logic (#93) --- faust/assignor/copartitioned_assignor.py | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/faust/assignor/copartitioned_assignor.py b/faust/assignor/copartitioned_assignor.py index 1c762ce85..b84434a92 100644 --- a/faust/assignor/copartitioned_assignor.py +++ b/faust/assignor/copartitioned_assignor.py @@ -73,7 +73,15 @@ def get_assignment(self) -> MutableMapping[str, CopartitionedAssignment]: # If there are no replicas required, try to ensure a balanced distribution with # at least one partition per client by removing partitions down to the # min_capacity for each client - capacity = self.min_capacity if self.replicas == 0 else self.max_capacity + capacity = ( + self.min_capacity + if ( + self.replicas == 0 + and self.num_partitions > 0 + and self.num_partitions >= self._num_clients + ) + else self.max_capacity + ) for copartitioned in self._client_assignments.values(): copartitioned.unassign_extras(capacity, self.replicas) self._assign(active=True)