Skip to content

Commit 9fca653

Browse files
committed
[Data] Export task USS distribution metrics
Signed-off-by: viiccwen <vicwen@apache.org>
1 parent ec368d7 commit 9fca653

7 files changed

Lines changed: 98 additions & 34 deletions

File tree

python/ray/dashboard/modules/metrics/dashboards/data_dashboard_panels.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -431,6 +431,21 @@
431431
stack=False,
432432
)
433433

434+
MAX_USS_PER_TASK_PANEL = Panel(
435+
id=91,
436+
title="Max Task USS Memory per Operator",
437+
description="Maximum unique set size (USS) memory usage in bytes among completed tasks for each operator.",
438+
unit="bytes",
439+
targets=[
440+
Target(
441+
expr='max(ray_data_max_uss_bytes_max{{{global_filters}, operator=~"$Operator"}}) by (dataset, operator)',
442+
legend="Max Task USS: {{dataset}}, {{operator}}",
443+
)
444+
],
445+
fill=0,
446+
stack=False,
447+
)
448+
434449
TASK_THROUGHPUT_BY_NODE_PANEL = Panel(
435450
id=46,
436451
title="Task Throughput (by Node)",
@@ -1479,6 +1494,7 @@
14791494
TASK_COMPLETION_TIME_WITHOUT_BACKPRESSURE_PANEL,
14801495
TASK_OUTPUT_BACKPRESSURE_TIME_PANEL,
14811496
TASK_SUBMISSION_BACKPRESSURE_PANEL,
1497+
MAX_USS_PER_TASK_PANEL,
14821498
TASK_THROUGHPUT_BY_NODE_PANEL,
14831499
TASKS_WITH_OUTPUT_PANEL,
14841500
SUBMITTED_TASKS_PANEL,

python/ray/data/_internal/execution/interfaces/op_runtime_metrics.py

Lines changed: 3 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ class MetricsType(Enum):
5454
Gauge = 1
5555
Histogram = 2
5656
Unsupported = 3
57+
Distribution = 4
5758

5859

5960
@dataclass(frozen=True)
@@ -899,21 +900,12 @@ def op_task_duration_stats(self) -> DistributionTracker:
899900
@metric_property(
900901
description="Distribution of max USS bytes across tasks.",
901902
metrics_group=MetricsGroup.TASKS,
902-
metrics_type=MetricsType.Unsupported,
903+
metrics_type=MetricsType.Distribution,
904+
metrics_args={"statistics": ("mean", "max")},
903905
)
904906
def max_uss_bytes(self) -> DistributionTracker:
905907
return self._max_uss_bytes
906908

907-
@metric_property(
908-
description="Average USS usage of tasks.",
909-
metrics_group=MetricsGroup.TASKS,
910-
)
911-
def average_max_uss_per_task(self) -> Optional[float]:
912-
"""Average max USS usage of tasks."""
913-
if self.max_uss_bytes.num_samples == 0:
914-
return None
915-
return self.max_uss_bytes.mean
916-
917909
@metric_property(
918910
description="Indicates if the operator is hanging.",
919911
metrics_group=MetricsGroup.MISC,

python/ray/data/_internal/issue_detection/detectors/high_memory_detector.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -81,19 +81,20 @@ def detect(self) -> List[Issue]:
8181
if not isinstance(op, MapOperator):
8282
continue
8383

84-
if op.metrics.average_max_uss_per_task is None:
84+
max_uss_bytes = op.metrics.max_uss_bytes
85+
if max_uss_bytes.num_samples == 0:
8586
continue
8687

8788
remote_args = op._get_dynamic_ray_remote_args()
8889
safe_memory_per_task = get_safe_default_logical_memory(remote_args)
8990

9091
if (
91-
op.metrics.average_max_uss_per_task > self._initial_memory_requests[op]
92-
and op.metrics.average_max_uss_per_task >= safe_memory_per_task
92+
max_uss_bytes.mean > self._initial_memory_requests[op]
93+
and max_uss_bytes.mean >= safe_memory_per_task
9394
):
9495
message = HIGH_MEMORY_PERIODIC_WARNING.format(
9596
op_name=op.name,
96-
memory_per_task=memory_string(op.metrics.average_max_uss_per_task),
97+
memory_per_task=memory_string(max_uss_bytes.mean),
9798
initial_memory_request=memory_string(
9899
self._initial_memory_requests[op]
99100
),

python/ray/data/_internal/stats.py

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -681,7 +681,7 @@ def __init__(self, max_stats=1000):
681681

682682
def _create_prometheus_metrics_for_execution_metrics(
683683
self, metrics_group: MetricsGroup, tag_keys: Tuple[str, ...]
684-
) -> Dict[str, Metric]:
684+
) -> Dict[str, Union[Metric, Dict[str, Gauge]]]:
685685
metrics = {}
686686
for metric in OpRuntimeMetrics.get_metrics():
687687
if not metric.metrics_group == metrics_group:
@@ -709,6 +709,15 @@ def _create_prometheus_metrics_for_execution_metrics(
709709
description=metric_description,
710710
tag_keys=tag_keys,
711711
)
712+
elif metric.metrics_type == MetricsType.Distribution:
713+
metrics[metric.name] = {
714+
statistic: Gauge(
715+
f"{metric_name}_{statistic}",
716+
description=f"{metric_description} ({statistic})",
717+
tag_keys=tag_keys,
718+
)
719+
for statistic in metric.metrics_args["statistics"]
720+
}
712721
return metrics
713722

714723
def _create_prometheus_metrics_for_per_node_metrics(self) -> Dict[str, Gauge]:
@@ -731,14 +740,14 @@ def gen_dataset_id(self) -> str:
731740
def update_execution_metrics(
732741
self,
733742
dataset_tag: str,
734-
op_metrics: List[Dict[str, int | float]],
743+
op_metrics: List[Dict[str, Any]],
735744
operator_tags: List[str],
736745
state: Dict[str, Any],
737746
per_node_metrics: Optional[Dict[str, Dict[str, int | float]]] = None,
738747
):
739748
def _record(
740-
prom_metric: Metric,
741-
value: Union[int, float, List[int]],
749+
prom_metric: Union[Metric, Dict[str, Gauge]],
750+
value: Any,
742751
tags: Dict[str, str] = None,
743752
):
744753
if isinstance(prom_metric, Gauge):
@@ -748,6 +757,13 @@ def _record(
748757
elif isinstance(prom_metric, Histogram):
749758
if isinstance(value, RuntimeMetricsHistogram):
750759
value.export_to(prom_metric, tags)
760+
elif isinstance(prom_metric, dict) and isinstance(value, dict):
761+
if value.get("num_samples", 0) == 0:
762+
return
763+
for statistic, gauge in prom_metric.items():
764+
statistic_value = value.get(statistic)
765+
if statistic_value is not None:
766+
gauge.set(statistic_value, tags)
751767

752768
for stats, operator_tag in zip(op_metrics, operator_tags):
753769
tags = self._create_tags(dataset_tag, operator_tag)

python/ray/data/tests/test_issue_detection.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -226,7 +226,9 @@ def test_high_memory_detection(
226226
data_context=ctx,
227227
ray_remote_args={"memory": configured_memory},
228228
)
229-
map_operator._metrics = MagicMock(average_max_uss_per_task=actual_memory)
229+
map_operator._metrics = MagicMock()
230+
map_operator._metrics.max_uss_bytes.num_samples = 1
231+
map_operator._metrics.max_uss_bytes.mean = actual_memory
230232
topology = {input_data_buffer: MagicMock(), map_operator: MagicMock()}
231233

232234
operators = list(topology.keys())

python/ray/data/tests/test_op_runtime_metrics.py

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,13 @@
1717
from ray.data.context import DataContext
1818

1919

20-
def test_average_max_uss_per_task():
20+
def test_max_uss_bytes_distribution():
2121
op = MagicMock()
2222
op.data_context.enable_get_object_locations_for_metrics = False
2323
metrics = OpRuntimeMetrics(op)
24-
assert metrics.average_max_uss_per_task is None
24+
assert metrics.max_uss_bytes.num_samples == 0
25+
assert metrics.max_uss_bytes.mean == 0
26+
assert metrics.max_uss_bytes.max is None
2527

2628
input_bundle = RefBundle([], owns_blocks=False, schema=None)
2729

@@ -33,7 +35,9 @@ def test_average_max_uss_per_task():
3335
TaskExecWorkerStats(task_wall_time_s=1.0, max_uss_bytes=100),
3436
TaskExecDriverStats(task_output_backpressure_s=0),
3537
)
36-
assert metrics.average_max_uss_per_task == 100
38+
assert metrics.max_uss_bytes.num_samples == 1
39+
assert metrics.max_uss_bytes.mean == 100
40+
assert metrics.max_uss_bytes.max == 100
3741

3842
# Submit and finish second task with USS of 300 bytes.
3943
metrics.on_task_submitted(1, input_bundle)
@@ -43,7 +47,9 @@ def test_average_max_uss_per_task():
4347
TaskExecWorkerStats(task_wall_time_s=1.0, max_uss_bytes=300),
4448
TaskExecDriverStats(task_output_backpressure_s=0),
4549
)
46-
assert metrics.average_max_uss_per_task == 200 # (100 + 300) / 2
50+
assert metrics.max_uss_bytes.num_samples == 2
51+
assert metrics.max_uss_bytes.mean == 200 # (100 + 300) / 2
52+
assert metrics.max_uss_bytes.max == 300
4753

4854

4955
def test_task_completion_time_histogram():

python/ray/data/tests/test_stats.py

Lines changed: 41 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@
5656
from ray.data.context import DataContext
5757
from ray.data.tests.util import column_udf
5858
from ray.tests.conftest import * # noqa
59+
from ray.util.metrics import Gauge
5960

6061

6162
@dataclass(frozen=True)
@@ -353,7 +354,6 @@ def gen_expected_metrics(
353354
"'average_rows_outputs_per_task': N",
354355
"'op_task_duration_stats': {'num_samples': N, 'mean': N, 'variance': N, 'min': N, 'max': N, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P}",
355356
"'max_uss_bytes': H",
356-
"'average_max_uss_per_task': H",
357357
"'num_inputs_received': N",
358358
"'num_row_inputs_received': N",
359359
"'bytes_inputs_received': N",
@@ -443,7 +443,6 @@ def gen_expected_metrics(
443443
"'average_rows_outputs_per_task': None",
444444
"'op_task_duration_stats': {'num_samples': Z, 'mean': Z, 'variance': Z, 'min': None, 'max': None, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P}",
445445
"'max_uss_bytes': H",
446-
"'average_max_uss_per_task': H",
447446
"'num_inputs_received': N",
448447
"'num_row_inputs_received': N",
449448
"'bytes_inputs_received': N",
@@ -635,11 +634,6 @@ def canonicalize(
635634
# Replace tabs with spaces.
636635
canonicalized_stats = re.sub("\t", " ", canonicalized_stats)
637636

638-
canonicalized_stats = re.sub(
639-
r"(average_max_uss_per_task:|'average_max_uss_per_task':) (?:N|Z|None)\b",
640-
r"\g<1> H",
641-
canonicalized_stats,
642-
)
643637
# Percentile values in DistributionTracker dicts can be None (when datasketches
644638
# is not installed) or a number (canonicalized to N). Normalize to P.
645639
canonicalized_stats = re.sub(
@@ -932,7 +926,6 @@ def test_dataset__repr__(ray_start_regular_shared, restore_data_context):
932926
" average_rows_outputs_per_task: N,\n"
933927
" op_task_duration_stats: {'num_samples': N, 'mean': N, 'variance': N, 'min': N, 'max': N, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P},\n"
934928
" max_uss_bytes: H,\n"
935-
" average_max_uss_per_task: H,\n"
936929
" num_inputs_received: N,\n"
937930
" num_row_inputs_received: N,\n"
938931
" bytes_inputs_received: N,\n"
@@ -1097,7 +1090,6 @@ def check_stats():
10971090
" average_rows_outputs_per_task: N,\n"
10981091
" op_task_duration_stats: {'num_samples': N, 'mean': N, 'variance': N, 'min': N, 'max': N, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P},\n"
10991092
" max_uss_bytes: H,\n"
1100-
" average_max_uss_per_task: H,\n"
11011093
" num_inputs_received: N,\n"
11021094
" num_row_inputs_received: N,\n"
11031095
" bytes_inputs_received: N,\n"
@@ -1215,7 +1207,6 @@ def check_stats():
12151207
" average_rows_outputs_per_task: N,\n"
12161208
" op_task_duration_stats: {'num_samples': N, 'mean': N, 'variance': N, 'min': N, 'max': N, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P, 'pN': P},\n"
12171209
" max_uss_bytes: H,\n"
1218-
" average_max_uss_per_task: H,\n"
12191210
" num_inputs_received: N,\n"
12201211
" num_row_inputs_received: N,\n"
12211212
" bytes_inputs_received: N,\n"
@@ -2017,6 +2008,46 @@ def test_stats_actor_iter_metrics():
20172008
assert update_fn.call_args_list[-1].args[2] is None
20182009

20192010

2011+
def test_stats_actor_exports_distribution_metrics():
2012+
actor = _StatsActor.__ray_metadata__.modified_class()
2013+
2014+
metrics = actor.execution_metrics_tasks["max_uss_bytes"]
2015+
assert set(metrics) == {"mean", "max"}
2016+
for statistic, metric in metrics.items():
2017+
assert isinstance(metric, Gauge)
2018+
assert metric.info["name"] == f"data_max_uss_bytes_{statistic}"
2019+
assert metric.info["tag_keys"] == ("dataset", "operator")
2020+
2021+
actor.update_dataset = MagicMock()
2022+
2023+
with (
2024+
patch.object(metrics["mean"], "set") as set_mean,
2025+
patch.object(metrics["max"], "set") as set_max,
2026+
):
2027+
actor.update_execution_metrics(
2028+
"dataset_1",
2029+
[{"max_uss_bytes": {"num_samples": 0, "mean": 0, "max": None}}],
2030+
["MapBatches_1"],
2031+
{},
2032+
)
2033+
set_mean.assert_not_called()
2034+
set_max.assert_not_called()
2035+
2036+
actor.update_execution_metrics(
2037+
"dataset_1",
2038+
[{"max_uss_bytes": {"num_samples": 2, "mean": 200, "max": 300}}],
2039+
["MapBatches_1"],
2040+
{},
2041+
)
2042+
2043+
set_mean.assert_called_once_with(
2044+
200, {"dataset": "dataset_1", "operator": "MapBatches_1"}
2045+
)
2046+
set_max.assert_called_once_with(
2047+
300, {"dataset": "dataset_1", "operator": "MapBatches_1"}
2048+
)
2049+
2050+
20202051
@pytest.mark.parametrize(
20212052
"split_index_arg, expected_split_label",
20222053
[

0 commit comments

Comments
 (0)