Skip to content

Commit 36d2607

Browse files
committed
benchmarking: add memory working set to server telemetry
1 parent c52b747 commit 36d2607

3 files changed

Lines changed: 216 additions & 14 deletions

File tree

‎benchmarking/README.md‎

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,7 @@ Three flags control the optional post-run measurements described in
9292
Defaults to the in-cluster service installed by
9393
[Optional: Prometheus + Grafana](#optional-prometheus--grafana).
9494
* `--atelet-lag-s`: how long to wait after the run before reading the
95-
atelet's snapshot metrics. Defaults to 70.
95+
atelet's metrics. Defaults to 70.
9696

9797
Test-specific flags are appended to the same command; see the sections below.
9898

@@ -456,6 +456,25 @@ actually did, independent of what the load generator reported.
456456
sample with no series counts as 0, since the atelet stops exporting when a
457457
node has no running actor; if the metric never appeared during the run,
458458
every field is `null`.
459+
* `working_set.node`, `ateom`, `atelet`: memory working set in GiB from
460+
cAdvisor on the nodes hosting a worker pod, per worker pod (every actor on it
461+
included), and per atelet on those nodes. Each has a percentile `summary`
462+
(plus `max`) over the steady-state window, `count` of series seen in it, and
463+
a `timeseries` every 10s over the whole run with the sum (`total_gb`) and
464+
largest (`max_gb`). A spike shorter than 10s can be missed.
465+
* `working_set.actor`, `per_actor`: the same shape for the actors' own working
466+
set as the atelet reports it (`ate_actor_stats_memory_working_set_bytes`),
467+
per atelet with every template summed, and the average per measured actor
468+
(`ate_actor_stats_sampled_actors`) on each atelet; `count` is atelets. As
469+
`per_actor` is an average per atelet, its `max` is the highest atelet
470+
average, not the largest actor, and its `total_gb` has no meaning. A gVisor
471+
actor is read from its sandbox cgroup on the host, a microVM actor from
472+
inside the guest. As with `active_actors`, suspended actors fall out, but
473+
through the telemetry-meter a series that stops being exported keeps its
474+
last value for up to 5m. Missing samples are skipped, not counted as 0. The
475+
atelet samples every `--actor-stats-poll-interval` (default 1m), so these
476+
move in steps of a minute or more. `actor` matches one worker pod's actors
477+
only when the node runs one worker pod.
459478

460479
Every distribution reports p50, p90, p95 and p99 over the steady-state window.
461480
The steady-state window runs from the first to the last Locust sample at 90% or
@@ -469,8 +488,9 @@ every 10s.
469488
The atelet exports before Prometheus scrapes it, so the harvest waits
470489
`--atelet-lag-s` seconds (default 70, enough for the OTel SDK's 60s default
471490
export and a 10s scrape)
472-
and reads the snapshot and `active_actors` windows half that late. The window
473-
ends at the last full-load sample, so teardown suspends are left out.
491+
and reads the snapshot, `active_actors` and actor working set windows half
492+
that late. The window ends at the last full-load sample, so teardown suspends
493+
are left out.
474494

475495
`metadata.start_ts` and `end_ts` bound the whole run, which the packing
476496
`timeseries` covers. `steady_start_ts` and `steady_end_ts` bound the

‎benchmarking/locust/server_telemetry.py‎

Lines changed: 135 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
1717
Queries Prometheus over [T_start, T_end] and the steady-state window [T_steady, T_end]
1818
to capture cluster packing, node and pod PSI stalls, snapshot sizes, latencies
19-
and throughput, and running actors.
19+
and throughput, running actors and memory working set.
2020
"""
2121

2222
import csv
@@ -547,6 +547,129 @@ def _count_active_actors(
547547
}
548548

549549

550+
# Live worker pods: the rate drops a deleted pod's stale series in ~1 min, not 5.
551+
WORKER_PODS = (
552+
"rate(container_cpu_usage_seconds_total"
553+
'{namespace="benchmark-workloads", container="ateom"}[1m]) > 0'
554+
)
555+
556+
# Only nodes/pods with a live worker; max, not sum: a restart briefly has 2 series.
557+
WORKING_SET_QUERIES = {
558+
"node": (
559+
'max by (instance) (container_memory_working_set_bytes{container="node"})'
560+
f" and on (instance) count by (instance) ({WORKER_PODS})"
561+
),
562+
"ateom": (
563+
"max by (pod) (container_memory_working_set_bytes"
564+
'{namespace="benchmark-workloads", container="ateom"})'
565+
f" and on (pod) ({WORKER_PODS})"
566+
),
567+
"atelet": (
568+
"max by (pod, instance) (container_memory_working_set_bytes"
569+
'{namespace="ate-system", container="atelet"})'
570+
f" and on (instance) count by (instance) ({WORKER_PODS})"
571+
),
572+
}
573+
574+
# Running actors per atelet, all templates summed; `instance` when scraped
575+
# directly, `exported_instance` via a collector. Not node names, so no worker
576+
# filter: the atelet stops exporting on a node with no running actor.
577+
ACTOR_WORKING_SET = (
578+
"sum by (instance, exported_instance) (ate_actor_stats_memory_working_set_bytes)"
579+
)
580+
ACTOR_COUNT = "sum by (instance, exported_instance) (ate_actor_stats_sampled_actors)"
581+
WORKING_SET_KEYS = (*WORKING_SET_QUERIES, "actor", "per_actor")
582+
583+
BYTES_PER_GIB = 2**30
584+
585+
586+
def _read_series(
587+
prom_url: str, query: str, start_ts: int, end_ts: int, shift: int = 0
588+
) -> dict[tuple, dict[int, float]]:
589+
"""Samples per series every 10s, read `shift` late and moved back by it."""
590+
res = query_prometheus_range(
591+
prom_url, query, start_ts + shift, end_ts + shift, step="10s"
592+
)
593+
out: dict[tuple, dict[int, float]] = {}
594+
for series in res:
595+
samples = out.setdefault(tuple(sorted(series.get("metric", {}).items())), {})
596+
for pt in series.get("values", []):
597+
try:
598+
t, val = int(pt[0]) - shift, float(pt[1])
599+
except (ValueError, IndexError, TypeError):
600+
continue
601+
if not math.isnan(val) and not math.isinf(val):
602+
samples[t] = val
603+
return out
604+
605+
606+
def _summarize_gb(
607+
series: dict[tuple, dict[int, float]],
608+
steady_start_ts: int,
609+
steady_end_ts: int,
610+
) -> dict[str, Any]:
611+
"""Steady percentiles, series seen in steady, and a 10s sum/max timeseries."""
612+
by_ts: dict[int, list[float]] = {}
613+
steady_gb: list[float] = []
614+
steady_series = 0
615+
for samples in series.values():
616+
in_steady = False
617+
for t, val in samples.items():
618+
gb = val / BYTES_PER_GIB
619+
by_ts.setdefault(t, []).append(gb)
620+
if steady_start_ts <= t <= steady_end_ts:
621+
steady_gb.append(gb)
622+
in_steady = True
623+
steady_series += in_steady
624+
return {
625+
"summary": compute_percentiles(steady_gb),
626+
"count": steady_series if by_ts else None,
627+
"timeseries": [
628+
{
629+
"timestamp": t,
630+
"total_gb": round(sum(v), 4),
631+
"max_gb": round(max(v), 4),
632+
}
633+
for t, v in sorted(by_ts.items())
634+
],
635+
}
636+
637+
638+
def _measure_working_set(
639+
prom_url: str,
640+
start_ts: int,
641+
end_ts: int,
642+
steady_start_ts: int,
643+
steady_end_ts: int,
644+
lag_s: int = 0,
645+
) -> dict[str, Any]:
646+
"""Working set GiB per node, ateom, atelet and actor; null, not 0, with no series.
647+
648+
`ateom` is the whole worker pod, every actor on it included. `actor` is the
649+
atelet's sum over the actors it measured, per atelet; `per_actor` divides it
650+
by their count. The atelet samples about once a minute and exports late, so
651+
both are read `lag_s // 2` late and moved back onto the run's 10s grid.
652+
"""
653+
def summarize(series: dict[tuple, dict[int, float]]) -> dict[str, Any]:
654+
return _summarize_gb(series, steady_start_ts, steady_end_ts)
655+
656+
out: dict[str, Any] = {
657+
key: summarize(_read_series(prom_url, query, start_ts, end_ts))
658+
for key, query in WORKING_SET_QUERIES.items()
659+
}
660+
shift = lag_s // 2
661+
actor = _read_series(prom_url, ACTOR_WORKING_SET, start_ts, end_ts, shift)
662+
count = _read_series(prom_url, ACTOR_COUNT, start_ts, end_ts, shift)
663+
out["actor"] = summarize(actor)
664+
# Same atelet and timestamp; a missing or zero count is skipped.
665+
out["per_actor"] = summarize({
666+
key: {t: v / count[key][t] for t, v in samples.items()
667+
if count.get(key, {}).get(t)}
668+
for key, samples in actor.items()
669+
})
670+
return out
671+
672+
550673
def harvest_server_telemetry(
551674
prom_url: str,
552675
start_ts: int,
@@ -557,9 +680,9 @@ def harvest_server_telemetry(
557680
) -> dict[str, Any]:
558681
"""Harvests the ground truth metric streams from Prometheus.
559682
560-
Only the atelet-exported snapshot and active actor blocks are read late, by
561-
up to `lag_s`; packing and PSI are scraped directly and read over the run
562-
window as is.
683+
Only the atelet-exported snapshot, active actor and actor working set
684+
blocks are read late, by up to `lag_s`; packing, PSI and the cAdvisor
685+
working set are scraped directly and read over the run window as is.
563686
"""
564687
steady_end = end_ts if steady_end_ts is None else steady_end_ts
565688
return {
@@ -578,6 +701,9 @@ def harvest_server_telemetry(
578701
"active_actors": _count_active_actors(
579702
prom_url, start_ts, end_ts, steady_start_ts, steady_end, lag_s
580703
),
704+
"working_set": _measure_working_set(
705+
prom_url, start_ts, end_ts, steady_start_ts, steady_end, lag_s
706+
),
581707
}
582708

583709

@@ -674,6 +800,11 @@ def log(msg: str) -> None:
674800
active = telemetry.get("active_actors", {}).get("summary", {})
675801
for p in ps:
676802
measurements[f"active_actors_{p}"] = active.get(p)
803+
working_set = telemetry.get("working_set", {})
804+
for block in WORKING_SET_KEYS:
805+
ws = working_set.get(block, {}).get("summary", {})
806+
for p in (*ps, "max"):
807+
measurements[f"working_set_{block}_{p}_gb"] = ws.get(p)
677808

678809
jsonl_row = {
679810
"timestamp": data_ts,

‎benchmarking/locust/unit_tests/test_server_telemetry.py‎

Lines changed: 58 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -40,8 +40,8 @@
4040
"end_ts": 105, "steady_start_ts": 100}
4141

4242

43-
def ranges(packing, node=(), pod=(), active=()):
44-
"""A query_prometheus_range fake answering packing, PSI and active actors.
43+
def ranges(packing, node=(), pod=(), active=(), ws=(), actor_ws=()):
44+
"""A query_prometheus_range fake: packing, PSI, active actors, working set.
4545
4646
Pod PSI only answers the pod-slice selector, so a query that loses the id
4747
regex (and would also sum the pause container) gets nothing back. Node PSI
@@ -56,6 +56,11 @@ def fake(_url, query, *_args, **_kwargs):
5656
return packing
5757
if "ate_actor_stats_sampled_actors" in query:
5858
return list(active)
59+
if "ate_actor_stats_memory_working_set_bytes" in query:
60+
return list(actor_ws)
61+
# Before the PSI branch: the node working set query has container="node".
62+
if "container_memory_working_set_bytes" in query:
63+
return list(ws)
5964
if 'container="node"' in query:
6065
return list(node) if worker_nodes in query and pod_slice in query else []
6166
return list(pod) if pod_slice in query else []
@@ -244,7 +249,11 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
244249
"values": [[100, "0.5"], [105, "0.5"]]} for p in "ab"]
245250
mock_range.side_effect = ranges(
246251
[partial, full, idle], node=quiet, pod=pods,
247-
active=[{"metric": {"instance": "atelet-a"}, "values": [[135, "3"]]}])
252+
active=[{"metric": {"instance": "atelet-a"}, "values": [[135, "3"]]}],
253+
ws=[{"metric": {"instance": "node-a"},
254+
"values": [[100, str(2**30)], [105, str(2**30)]]}],
255+
actor_ws=[{"metric": {"instance": "atelet-a"},
256+
"values": [[135, str(6 * 2**30)]]}])
248257
# The default 70s lag reads the snapshot window at 135 and 140.
249258
mock_instant.side_effect = snapshot_prom(
250259
"11.5",
@@ -278,14 +287,14 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
278287

279288
sleep.assert_called_once_with(65) # until end 105 + 70, from 110
280289
self.assertEqual(summary["metadata"]["atelet_lag_s"], 70)
281-
# Packing and PSI are not shifted; active actors are, like snapshots.
282-
windows = {("sampled_actors" in c.args[1], c.args[2:4])
290+
# Packing, PSI and cAdvisor are not shifted; atelet stats are, like snapshots.
291+
windows = {("ate_actor_stats" in c.args[1], c.args[2:4])
283292
for c in mock_range.call_args_list}
284293
self.assertEqual(windows, {(False, (100, 105)), (True, (135, 140))})
285294
self.assertEqual(
286295
set(summary),
287296
{"metadata", "cluster_packing", "node_psi", "pod_psi", "snapshots",
288-
"active_actors"},
297+
"active_actors", "working_set"},
289298
)
290299
packing = summary["cluster_packing"]
291300
self.assertEqual(packing["summary"]["p50"], 0.8) # 3 + 1 busy / 5 workers
@@ -305,7 +314,7 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
305314
self.assertEqual(row["metric"], "server_summary")
306315
m = row["measurements"]
307316
# Every source answered, so a key wired to a wrong name would read None.
308-
self.assertEqual(len(m), 49)
317+
self.assertEqual(len(m), 74)
309318
self.assertEqual([k for k, v in m.items() if v is None], [])
310319
# Every value a string, so one row's types match every other row's.
311320
self.assertTrue(all(isinstance(v, str) for v in m.values()))
@@ -317,6 +326,11 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
317326
self.assertEqual(m["restore_mean_s"], "0.5")
318327
self.assertEqual(m["checkpoint_mean_s"], "2.0")
319328
self.assertEqual(m["active_actors_p99"], "3.0")
329+
# Read at 135, moved back to 100; 6 GiB over the 3 actors on atelet-a.
330+
actor = summary["working_set"]["actor"]
331+
self.assertEqual(actor["timeseries"][0]["timestamp"], 100)
332+
self.assertEqual(m["working_set_actor_p50_gb"], "6.0")
333+
self.assertEqual(m["working_set_per_actor_p50_gb"], "2.0")
320334

321335
@mock.patch("server_telemetry.query_prometheus_range")
322336
@mock.patch("server_telemetry.query_prometheus_instant")
@@ -466,6 +480,43 @@ def test_steady_window_bounds_packing_and_psi(self, mock_instant, mock_range):
466480
self.assertEqual(out["node_psi"]["cpu_stall_pct"]["max"], 0.5)
467481
self.assertEqual(out["pod_psi"]["cpu_stall_pct"]["max"], 0.5)
468482

483+
@mock.patch("server_telemetry.query_prometheus_range")
484+
def test_working_set_gib_and_steady_window(self, mock_range):
485+
gib = 2**30
486+
mock_range.side_effect = ranges([], ws=[
487+
{"metric": {"instance": "node-a"},
488+
"values": [[90, str(8 * gib)], [100, str(gib)], [110, str(3 * gib)]]},
489+
{"metric": {"instance": "node-b"}, "values": [[110, str(gib)]]},
490+
])
491+
out = server_telemetry._measure_working_set("http://p", 90, 110, 100, 110)
492+
node = out["node"]
493+
# The 8 GiB sample at 90 is ramp-up: in the timeseries, not the summary.
494+
self.assertEqual((node["summary"]["p50"], node["summary"]["max"]), (1.0, 3.0))
495+
self.assertEqual(node["count"], 2)
496+
self.assertEqual(node["timeseries"], [
497+
{"timestamp": 90, "total_gb": 8.0, "max_gb": 8.0},
498+
{"timestamp": 100, "total_gb": 1.0, "max_gb": 1.0},
499+
{"timestamp": 110, "total_gb": 4.0, "max_gb": 3.0},
500+
])
501+
self.assertEqual(set(out), {"node", "ateom", "atelet", "actor", "per_actor"})
502+
503+
@mock.patch("server_telemetry.query_prometheus_range")
504+
def test_working_set_null_not_zero(self, mock_range):
505+
mock_range.side_effect = ranges([])
506+
out = server_telemetry._measure_working_set("http://p", 90, 110, 100, 110)
507+
for block in out.values():
508+
self.assertEqual(block, {"summary": NO_PERCENTILES, "count": None,
509+
"timeseries": []})
510+
511+
def test_working_set_queries_keep_live_workers(self):
512+
q = server_telemetry.WORKING_SET_QUERIES
513+
live = "rate(container_cpu_usage_seconds_total"
514+
for block in q.values():
515+
self.assertTrue(block.startswith("max by (")) # a restart has 2 series
516+
self.assertIn(live, block)
517+
self.assertIn('container="node"', q["node"])
518+
self.assertIn("and on (pod)", q["ateom"])
519+
469520
@mock.patch("server_telemetry.query_prometheus_instant")
470521
def test_failed_delta_read_is_unknown(self, mock_instant):
471522
# A failed delta query is not taken as a series born in the window.

0 commit comments

Comments
 (0)