[Autoscaler] Placement group autoscaling (#11243)

This commit is contained in:
Alex Wu
2020-10-14 13:11:46 -07:00
committed by GitHub
parent aefcf901d3
commit 7466ce82df
5 changed files with 455 additions and 58 deletions
@@ -17,6 +17,8 @@ from ray.autoscaler._private.commands import get_or_create_head_node
from ray.autoscaler._private.resource_demand_scheduler import \
_utilization_score, _add_min_workers_nodes, \
get_bin_pack_residual, get_nodes_for, ResourceDemandScheduler
from ray.gcs_utils import PlacementGroupTableData
from ray.core.generated.common_pb2 import Bundle, PlacementStrategy
from ray.autoscaler.tags import TAG_RAY_USER_NODE_TYPE, TAG_RAY_NODE_KIND, \
NODE_KIND_WORKER, TAG_RAY_NODE_STATUS, \
STATUS_UP_TO_DATE, STATUS_UNINITIALIZED
@@ -94,15 +96,44 @@ def test_util_score():
def test_bin_pack():
assert get_bin_pack_residual([], [{"GPU": 2}, {"GPU": 2}]) == \
assert get_bin_pack_residual([], [{"GPU": 2}, {"GPU": 2}])[0] == \
[{"GPU": 2}, {"GPU": 2}]
assert get_bin_pack_residual([{"GPU": 2}], [{"GPU": 2}, {"GPU": 2}]) == \
[{"GPU": 2}]
assert get_bin_pack_residual([{"GPU": 4}], [{"GPU": 2}, {"GPU": 2}]) == []
assert get_bin_pack_residual([{"GPU": 2}], [{"GPU": 2}, {"GPU": 2}])[0] \
== [{"GPU": 2}]
assert get_bin_pack_residual([{
"GPU": 4
}], [{
"GPU": 2
}, {
"GPU": 2
}])[0] == []
arg = [{"GPU": 2}, {"GPU": 2, "CPU": 2}]
assert get_bin_pack_residual(arg, [{"GPU": 2}, {"GPU": 2}]) == []
assert get_bin_pack_residual(arg, [{"GPU": 2}, {"GPU": 2}])[0] == []
arg = [{"CPU": 2}, {"GPU": 2}]
assert get_bin_pack_residual(arg, [{"GPU": 2}, {"GPU": 2}]) == [{"GPU": 2}]
assert get_bin_pack_residual(arg, [{
"GPU": 2
}, {
"GPU": 2
}])[0] == [{
"GPU": 2
}]
arg = [{"GPU": 3}]
assert get_bin_pack_residual(
arg, [{
"GPU": 1
}, {
"GPU": 1
}], strict_spread=False)[0] == []
assert get_bin_pack_residual(
arg, [{
"GPU": 1
}, {
"GPU": 1
}], strict_spread=True) == ([{
"GPU": 1
}], [{
"GPU": 2
}])
def test_get_nodes_packing_heuristic():
@@ -135,6 +166,18 @@ def test_get_nodes_packing_heuristic():
assert get_nodes_for(
TYPES_A, {}, 9999, ([{"GPU": 1}] * 8) + ([{"CPU": 1}] * 64)) == \
{"m4.16xlarge": 1, "p2.8xlarge": 1}
assert get_nodes_for(
TYPES_A, {}, 9999, [{
"GPU": 1
}] * 8, strict_spread=False) == {
"p2.8xlarge": 1
}
assert get_nodes_for(
TYPES_A, {}, 9999, [{
"GPU": 1
}] * 8, strict_spread=True) == {
"p2.xlarge": 8
}
def test_get_nodes_respects_max_limit():
@@ -248,7 +291,7 @@ def test_get_nodes_to_launch_with_min_workers():
to_launch = scheduler.get_nodes_to_launch(nodes, {}, [{
"GPU": 8
}], utilizations)
}], utilizations, [])
assert to_launch == {"p2.8xlarge": 1}
@@ -270,7 +313,7 @@ def test_get_nodes_to_launch_with_min_workers_and_bin_packing():
# requires 2 p2.8xls (only 2 are in cluster/pending) and 1 p2.xlarge
demands = [{"GPU": 8}] * (len(utilizations) + 1) + [{"GPU": 1}]
to_launch = scheduler.get_nodes_to_launch(nodes, pending_nodes, demands,
utilizations)
utilizations, [])
assert to_launch == {"p2.xlarge": 1}
# 3 min_workers of p2.8xlarge covers the 2 p2.8xlarge + 1 p2.xlarge demand.
@@ -279,7 +322,7 @@ def test_get_nodes_to_launch_with_min_workers_and_bin_packing():
new_types["p2.8xlarge"]["min_workers"] = 3
scheduler = ResourceDemandScheduler(provider, new_types, 10)
to_launch = scheduler.get_nodes_to_launch(nodes, pending_nodes, demands,
utilizations)
utilizations, [])
# Make sure it does not return [("p2.8xlarge", 1), ("p2.xlarge", 1)]
assert to_launch == {"p2.8xlarge": 1}
@@ -297,7 +340,7 @@ def test_get_nodes_to_launch_limits():
to_launch = scheduler.get_nodes_to_launch(nodes, {"p2.8xlarge": 1}, [{
"GPU": 8
}] * 2, utilizations)
}] * 2, utilizations, [])
assert to_launch == {}
@@ -317,11 +360,101 @@ def test_calculate_node_resources():
# requires 4 p2.8xls (only 3 are in cluster/pending)
demands = [{"GPU": 8}] * (len(utilizations) + 2)
to_launch = scheduler.get_nodes_to_launch(nodes, pending_nodes, demands,
utilizations)
utilizations, [])
assert to_launch == {"p2.8xlarge": 1}
class TestPlacementGroupScaling:
def test_strategies(self):
provider = MockProvider()
scheduler = ResourceDemandScheduler(provider, TYPES_A, 10)
provider.create_node({}, {TAG_RAY_USER_NODE_TYPE: "p2.8xlarge"}, 2)
# At this point our cluster has 2 p2.8xlarge instances (16 GPUs) and is
# fully idle.
nodes = provider.non_terminated_nodes({})
resource_demands = [{"GPU": 4}] * 2
pending_placement_groups = [
# Requires a new node (only uses 2 GPUs on it though).
PlacementGroupTableData(
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.STRICT_SPREAD,
bundles=[
Bundle(unit_resources={"GPU": 2}),
Bundle(unit_resources={"GPU": 2}),
Bundle(unit_resources={"GPU": 2})
]),
# Requires a new node (uses the whole node).
PlacementGroupTableData(
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.STRICT_PACK,
bundles=([Bundle(unit_resources={"GPU": 2})] * 4)),
# Fits across the machines that strict spread.
PlacementGroupTableData(
# runs on.
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.PACK,
bundles=([Bundle(unit_resources={"GPU": 2})] * 2)),
# Fits across the machines that strict spread.
PlacementGroupTableData(
# runs on.
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.SPREAD,
bundles=([Bundle(unit_resources={"GPU": 2})] * 2)),
]
to_launch = scheduler.get_nodes_to_launch(nodes, {}, resource_demands,
{}, pending_placement_groups)
assert to_launch == {"p2.8xlarge": 2}
def test_many_strict_spreads(self):
provider = MockProvider()
scheduler = ResourceDemandScheduler(provider, TYPES_A, 10)
provider.create_node({}, {TAG_RAY_USER_NODE_TYPE: "p2.8xlarge"}, 2)
# At this point our cluster has 2 p2.8xlarge instances (16 GPUs) and is
# fully idle.
nodes = provider.non_terminated_nodes({})
resource_demands = [{"GPU": 1}] * 6
pending_placement_groups = [
# Requires a new node (only uses 2 GPUs on it though).
PlacementGroupTableData(
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.STRICT_SPREAD,
bundles=[Bundle(unit_resources={"GPU": 2})] * 3),
]
# Each placement group will take up 2 GPUs per node, but the distinct
# placement groups should still reuse the same nodes.
pending_placement_groups = pending_placement_groups * 3
to_launch = scheduler.get_nodes_to_launch(nodes, {}, resource_demands,
{}, pending_placement_groups)
assert to_launch == {"p2.8xlarge": 1}
def test_packing(self):
provider = MockProvider()
scheduler = ResourceDemandScheduler(provider, TYPES_A, 10)
provider.create_node({}, {TAG_RAY_USER_NODE_TYPE: "p2.8xlarge"}, 1)
# At this point our cluster has 1 p2.8xlarge instances (8 GPUs) and is
# fully idle.
nodes = provider.non_terminated_nodes({})
resource_demands = [{"GPU": 1}] * 2
pending_placement_groups = [
PlacementGroupTableData(
state=PlacementGroupTableData.PENDING,
strategy=PlacementStrategy.STRICT_PACK,
bundles=[Bundle(unit_resources={"GPU": 2})] * 3),
]
# The 2 resource demand gpus should still be packed onto the same node
# as the 6 GPU placement group.
to_launch = scheduler.get_nodes_to_launch(nodes, {}, resource_demands,
{}, pending_placement_groups)
assert to_launch == {}
def test_get_concurrent_resource_demand_to_launch():
node_types = copy.deepcopy(TYPES_A)
node_types["p2.8xlarge"]["min_workers"] = 1
@@ -441,7 +574,7 @@ def test_get_nodes_to_launch_max_launch_concurrency():
scheduler = ResourceDemandScheduler(provider, new_types, 30)
to_launch = scheduler.get_nodes_to_launch([], {}, [], [])
to_launch = scheduler.get_nodes_to_launch([], {}, [], [], [])
# Respects min_workers despite concurrency limitation.
assert to_launch == {"p2.8xlarge": 4}
@@ -456,7 +589,7 @@ def test_get_nodes_to_launch_max_launch_concurrency():
# requires 41 p2.8xls (currently 1 pending, 1 launching, 0 running}
demands = [{"GPU": 8}] * (len(utilizations) + 40)
to_launch = scheduler.get_nodes_to_launch(nodes, launching_nodes, demands,
utilizations)
utilizations, [])
# Enforces max launch to 5 when < 5 running. 2 are pending/launching.
assert to_launch == {"p2.8xlarge": 3}
@@ -471,7 +604,7 @@ def test_get_nodes_to_launch_max_launch_concurrency():
# requires 17 p2.8xls (currently 1 pending, 1 launching, 8 running}
demands = [{"GPU": 8}] * (len(utilizations) + 15)
to_launch = scheduler.get_nodes_to_launch(nodes, launching_nodes, demands,
utilizations)
utilizations, [])
# We are allowed to launch up to 8 more since 8 are running.
# We already have 2 pending/launching, so only 6 remain.
assert to_launch == {"p2.8xlarge": 6}
@@ -496,6 +629,25 @@ class LoadMetricsTest(unittest.TestCase):
"GPU": 1
}])
def testPlacementGroupLoad(self):
lm = LoadMetrics()
pending_placement_groups = [
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.PACK,
bundles=([Bundle(unit_resources={"GPU": 2})] * 2)),
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.SPREAD,
bundles=([Bundle(unit_resources={"GPU": 2})] * 2)),
]
lm.update(
"1.1.1.1", {},
True, {},
True, {},
pending_placement_groups=pending_placement_groups)
assert lm.get_pending_placement_groups() == pending_placement_groups
class AutoscalingTest(unittest.TestCase):
def setUp(self):
@@ -577,6 +729,75 @@ class AutoscalingTest(unittest.TestCase):
autoscaler.update()
self.waitForNodes(2)
def testPlacementGroup(self):
# Note this is mostly an integration test. See
# testPlacementGroupScaling for more comprehensive tests.
config = copy.deepcopy(MULTI_WORKER_CLUSTER)
config["min_workers"] = 0
config["max_workers"] = 999
config_path = self.write_config(config)
self.provider = MockProvider()
runner = MockProcessRunner()
lm = LoadMetrics()
autoscaler = StandardAutoscaler(
config_path,
lm,
max_failures=0,
process_runner=runner,
update_interval_s=0)
self.provider.create_node({}, {
TAG_RAY_NODE_KIND: "head",
TAG_RAY_USER_NODE_TYPE: "m4.4xlarge"
}, 1)
head_ip = self.provider.non_terminated_node_ips({})[0]
assert len(self.provider.non_terminated_nodes({})) == 1
autoscaler.update()
self.waitForNodes(1)
pending_placement_groups = [
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.STRICT_SPREAD,
bundles=[Bundle(unit_resources={"GPU": 2})] * 3),
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.PACK,
bundles=([Bundle(unit_resources={"GPU": 2})] * 5)),
]
# Since placement groups are implemented with custom resources, this is
# an example of the accompanying resource demands. Note the resource
# demand autoscaler will be unable to fulfill these demands, but we
# should still handle the other infeasible/waiting bundles.
placement_group_resource_demands = [{
"GPU_group_0_6c2506ac733bc37496295b02c4fad446": 0.0101,
"GPU_group_6c2506ac733bc37496295b02c4fad446": 0.0101
}]
lm.update(
head_ip, {"CPU": 16},
True, {"CPU": 16},
False, {},
infeasible_bundles=placement_group_resource_demands,
waiting_bundles=[{
"GPU": 8
}],
pending_placement_groups=pending_placement_groups)
autoscaler.update()
self.waitForNodes(5)
for i in range(1, 5):
assert self.provider.mock_nodes[i].node_type == "p2.8xlarge"
pending_placement_groups = [
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.STRICT_PACK,
bundles=([Bundle(unit_resources={"GPU": 2})] * 4)),
PlacementGroupTableData(
state=PlacementGroupTableData.RESCHEDULING,
strategy=PlacementStrategy.SPREAD,
bundles=([Bundle(unit_resources={"GPU": 2})] * 2)),
]
def testScaleUpMinWorkers(self):
config = copy.deepcopy(MULTI_WORKER_CLUSTER)
config["min_workers"] = 2