Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions datasets_eval/latency/.gitignore
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
replay-warm.jsonl
replay-cold.jsonl
61 changes: 61 additions & 0 deletions datasets_eval/latency/agent-cold-ship.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
# Cold-ON edge agent for the Fig 7 cold-fallback latency arm.
#
# Derived from datasets_eval/multisketch/agent-ddsketch-coldon.yaml, but
# with a COMPLETE cold block (ship_endpoint + block_duration +
# external_labels) so the edge actually SHIPS cold ASAPFRG1 fragments to
# the gorilla-merger. The multisketch coldon config only set
# `cold: {enabled: true}` with NO ship_endpoint, which the asapedge
# processor treats as DRAIN-ONLY (no shipping) — that is why the
# original cold arm landed an empty MinIO (config.go: "Empty =>
# drain-only (no shipping)").
#
# The lossless raw `google_cluster_2019_cpu_rate` Sum family (tier: both)
# is the series that ships cold; the data-plane routes its queries to the
# gorilla cold tier (backend-storage-routing-coldon.yaml).
receivers:
otlp:
protocols:
grpc: {endpoint: 0.0.0.0:4317, max_recv_msg_size_mib: 4096}
http: {endpoint: 0.0.0.0:4318}
processors:
batch:
send_batch_size: 800
send_batch_max_size: 1500
memory_limiter: {check_interval: 1s, limit_mib: 8192, spike_limit_mib: 1024}
asap_edge:
shard_count: 12
window_duration: 60s
drop_original: true
max_series: 200000
delta_transmission: false
metrics:
- {metric: google_cluster_2019_cpu_rate, family: sum, tier: both}
- {metric: google_cluster_2019_memory_usage, family: sum, aggregate_by: [zone], tier: both}
- metric: google_cluster_2019_cpu_rate_q_ddsketch
family: "ddsketch"
tier: "both"
delta_transmission: true
cold:
enabled: true
ship_endpoint: http://gorilla-merger:10908/ingest/gorilla
block_duration: 60s
reorder_grace: 2s
external_labels:
# MUST be non-empty: Thanos rejects a block with empty external
# labels ("empty external labels are not allowed for Thanos
# block"); matches the merger's --external-labels cluster=asap-mvp.
cluster: asap-mvp
control_channel: {enabled: false}
exporters:
otlp/backend:
endpoint: data-plane:14317
tls: {insecure: true}
timeout: 120s
sending_queue: {enabled: true, num_consumers: 4, queue_size: 5000}
service:
pipelines:
metrics: {receivers: [otlp], processors: [memory_limiter, asap_edge, batch], exporters: [otlp/backend]}
telemetry:
metrics:
level: detailed
readers: [{pull: {exporter: {prometheus: {host: 0.0.0.0, port: 8890}}}}]
29 changes: 29 additions & 0 deletions datasets_eval/latency/backend-storage-routing-coldon.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
# Cold-arm storage-routing table (Fig 7 cold-fallback latency arm).
#
# Routes the lossless raw `google_cluster_2019_cpu_rate` series — the
# Sum family the cold-enabled agent ships to the gorilla cold tier
# (MinIO via the merger) — to `gorilla_object_store`, so its PromQL
# queries are answered by the ThanosQueryEngine (data_source =
# thanos_query). This is what makes the COLD/archive path observable
# end-to-end: an instant query against this metric is dispatched
# through EngineRouter to thanos-query, which resolves the answer over
# the gorilla-merger StoreAPI + thanos-store-gateway (S3/MinIO blocks).
#
# `double_write` routes the metric to BOTH warm sketches and the
# archive, with the cost-aware dispatcher picking per query. We use
# the single-target `gorilla_object_store` here to GUARANTEE every
# timed query lands on the cold path (no warm shortcut) — we are
# measuring the cold arm, so we force the cold engine.
#
# Used by datasets_eval/latency/stack-coldon.sh (cold-ON stack).

default: sketch_store

routes:
- metric: google_cluster_2019_cpu_rate
targets:
- backend: gorilla_object_store
# all shapes route to the cold archive — sum / count /
# quantile_over_time over the raw shipped samples are all
# computed by thanos-query over the merger StoreAPI.
applies_to_query_shape: [count, topk, rate_post_hoc, quantile, sum, last_over_time, other]
128 changes: 128 additions & 0 deletions datasets_eval/latency/cold_latency_replay.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
#!/usr/bin/env python3
"""Cold-arm query-latency replay (Fig 7 cold-fallback arm).

Mirrors deploy/mvp-singlenode/scripts/metricsql_replay.py (warm arm) but
pins the PromQL *evaluation timestamp* (`time=<at_time>`) to the instant
the cold workload was anchored at, so every instant query deterministically
intersects the cold-archived window in MinIO/Thanos. (The warm arm queried
at live `now` because its in-memory warm window sat at `now`; the cold
window sits at a fixed past instant — the ship takes ~one window to land —
so we pin the eval time to it. The pin changes only WHICH timestamp the
backend evaluates at; the per-query server-side latency it measures is
identical in kind to the warm arm.)

Fires the query mix at a fixed QPS against the data-plane query surface
(:9091/api/v1/query), captures per-query wall-clock latency + the served
`data_source` (asap_query=warm vs thanos_query=cold) + result vector, and
writes a JSONL identical in schema to the warm replay so
compute_latency.py reduces both arms the same way.

Usage:
cold_latency_replay.py --target http://127.0.0.1:9091 \
--queries queries-latency-cold.json --at-time <unix_seconds> \
--qps 15 --duration 40 --out replay-cold.jsonl
"""
from __future__ import annotations
import argparse, datetime as dt, json, sys, threading, time
import urllib.parse, urllib.request, urllib.error


def run_query(target: str, metricsql: str, at_time: float | None,
timeout_s: float = 10.0):
params = {"query": metricsql}
if at_time is not None:
params["time"] = f"{at_time:.3f}"
url = f"{target.rstrip('/')}/api/v1/query?" + urllib.parse.urlencode(params)
started = time.perf_counter()
try:
with urllib.request.urlopen(url, timeout=timeout_s) as resp:
code = resp.getcode()
body = resp.read().decode("utf-8", errors="replace")
except urllib.error.HTTPError as e:
return (time.perf_counter() - started) * 1000.0, {
"status": "http_error", "http_code": e.code, "result": None,
"result_type": None, "data_source": None, "error": str(e)}
except Exception as e:
return (time.perf_counter() - started) * 1000.0, {
"status": "timeout", "http_code": None, "result": None,
"result_type": None, "data_source": None, "error": str(e)}
duration_ms = (time.perf_counter() - started) * 1000.0
try:
parsed = json.loads(body)
except json.JSONDecodeError as e:
return duration_ms, {"status": "json_error", "http_code": code,
"result": None, "result_type": None,
"data_source": None, "error": str(e)}
data = parsed.get("data") or {}
data_source = None
for info in (parsed.get("infos") or []):
if isinstance(info, str) and info.startswith("data_source:"):
data_source = info.split(":", 1)[1].strip()
return duration_ms, {
"status": parsed.get("status", "unknown"), "http_code": code,
"result": data.get("result"), "result_type": data.get("resultType"),
"data_source": data_source}


def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--target", default="http://127.0.0.1:9091")
ap.add_argument("--queries", required=True)
ap.add_argument("--at-time", type=float, default=None,
help="Unix seconds to pin the PromQL eval timestamp to. "
"Omit to query live now (warm-style).")
ap.add_argument("--qps", type=float, default=15.0)
ap.add_argument("--duration", type=float, default=40.0)
ap.add_argument("--out", required=True)
args = ap.parse_args()

with open(args.queries) as f:
queries = json.load(f)
if not queries:
sys.exit("no queries")

interval = 1.0 / args.qps
deadline = time.time() + args.duration
lock = threading.Lock()
rows: list[dict] = []
qi = 0
n_sent = 0
t0 = time.time()
while time.time() < deadline:
q = queries[qi % len(queries)]
qi += 1
dur_ms, res = run_query(args.target, q["metricsql"], args.at_time)
rec = {
"ts": dt.datetime.now(dt.timezone.utc).isoformat(),
"query": q["metricsql"], "kind": q.get("kind", "other"),
"duration_ms": round(dur_ms, 4), "status": res["status"],
"http_code": res.get("http_code"),
"result_type": res.get("result_type"),
"result": res.get("result"),
"data_source": res.get("data_source"),
"n_result_series": len(res.get("result") or []),
}
rows.append(rec)
n_sent += 1
# fixed-rate pacing
next_at = t0 + n_sent * interval
sleep = next_at - time.time()
if sleep > 0:
time.sleep(sleep)

with open(args.out, "w") as f:
for r in rows:
f.write(json.dumps(r) + "\n")
# quick summary
from collections import Counter
ds = Counter(r["data_source"] for r in rows)
st = Counter(r["status"] for r in rows)
empt = sum(1 for r in rows if not r["result"])
print(f"replay: {len(rows)} queries -> {args.out}")
print(f" data_source: {dict(ds)}")
print(f" status: {dict(st)} empty_results: {empt}")
return 0


if __name__ == "__main__":
raise SystemExit(main())
99 changes: 77 additions & 22 deletions datasets_eval/latency/compute_latency.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,32 +41,51 @@ def pct(sorted_vals, p):
return sorted_vals[lo] + (sorted_vals[hi] - sorted_vals[lo]) * frac


def summarize(rows, label, require_warm=True):
def _n_result(r) -> int:
"""Result count from either the full replay schema (`result` vector) or
the slim per-query schema (`n_result_series`)."""
res = r.get("result")
if res is not None:
return len(res)
return int(r.get("n_result_series") or 0)


def summarize(rows, label, require_real=True, require_source=None):
"""Returns (summary_dict, {kind: [latencies]}, [all_latencies]).
Hard-fails if require_warm and any timed query was empty/errored."""
Hard-fails if require_real and any timed query was empty/errored, or if
require_source is set and any timed query was served by a different
data_source (guards that the arm's answers came from the intended tier)."""
bad = []
wrong_src = []
by_kind: dict[str, list[float]] = {}
alllat: list[float] = []
src_counter: dict[str, int] = {}
for r in rows:
res = r.get("result") or []
n = _n_result(r)
src = r.get("data_source")
src_counter[src] = src_counter.get(src, 0) + 1
ok = r.get("status") == "success" and bool(res)
if require_warm and not ok:
ok = r.get("status") == "success" and n > 0
if require_real and not ok:
bad.append({"query": r.get("query"), "status": r.get("status"),
"n_result": len(res), "data_source": src})
"n_result": n, "data_source": src})
continue
if require_source is not None and src != require_source:
wrong_src.append({"query": r.get("query"), "data_source": src})
continue
if not ok:
continue
lat = float(r["duration_ms"])
alllat.append(lat)
by_kind.setdefault(r["kind"], []).append(lat)

if require_warm and bad:
if require_real and bad:
sys.exit(f"[{label}] {len(bad)} timed queries returned no real result "
f"(empty/errored) — latency would be meaningless. First few: "
f"{bad[:3]}")
if require_source is not None and wrong_src:
sys.exit(f"[{label}] {len(wrong_src)} timed queries were NOT served by "
f"data_source={require_source} (wrong tier — this arm must "
f"measure that tier only). First few: {wrong_src[:3]}")

def stats(vals):
s = sorted(vals)
Expand Down Expand Up @@ -98,24 +117,41 @@ def cdf_xy(vals):
return s, ys


def load_any(path: str):
"""Load either a JSONL replay log (one object per line) or a JSON array
(the committed slim per_query_latency*.json)."""
text = Path(path).read_text().lstrip()
if text.startswith("["):
return json.loads(text)
return load(path)


def main():
ap = argparse.ArgumentParser()
ap.add_argument("--warm", required=True)
ap.add_argument("--cold", default=None)
ap.add_argument("--warm", required=True,
help="warm replay JSONL or committed per_query_latency.json")
ap.add_argument("--cold", default=None,
help="cold replay JSONL or per_query_latency_cold.json")
ap.add_argument("--out-json", required=True)
ap.add_argument("--out-png", required=True)
args = ap.parse_args()

warm_rows = load(args.warm)
warm_sum, warm_by_kind, warm_all = summarize(warm_rows, "warm", require_warm=True)
warm_rows = load_any(args.warm)
# Warm arm: every timed query must be a real warm answer (data_source=asap_query).
warm_sum, warm_by_kind, warm_all = summarize(
warm_rows, "warm", require_real=True, require_source="asap_query")

out = {"warm": warm_sum}

cold_all = None
cold_by_kind = None
if args.cold and Path(args.cold).exists():
cold_rows = load(args.cold)
# cold arm: don't hard-fail on empties (route may differ); report what landed
cold_sum, _, cold_all = summarize(cold_rows, "cold", require_warm=False)
cold_rows = load_any(args.cold)
# Cold arm: GUARD that every timed query landed real AND was served by
# the cold/archive engine (data_source=thanos_query) — we are measuring
# the cold-fallback tier, so a warm shortcut would invalidate the arm.
cold_sum, cold_by_kind, cold_all = summarize(
cold_rows, "cold", require_real=True, require_source="thanos_query")
out["cold"] = cold_sum

Path(args.out_json).write_text(json.dumps(out, indent=2) + "\n")
Expand All @@ -136,28 +172,47 @@ def main():
color=palette.get(kind, "#7f7f7f"), lw=1.4, ls="--", alpha=0.85)
if cold_all:
x, y = cdf_xy(cold_all)
ax.plot(x, y, label=f"cold-fallback ({len(cold_all)} q)",
color="#d62728", lw=2.0)

ax.plot(x, y, label=f"cold-fallback — all ({len(cold_all)} q)",
color="#d62728", lw=2.2)
cold_palette = {"quantile": "#8c564b", "sum": "#e377c2"}
for kind, vals in sorted((cold_by_kind or {}).items()):
x, y = cdf_xy(vals)
ax.plot(x, y, label=f"cold — {kind} ({len(vals)} q)",
color=cold_palette.get(kind, "#d62728"), lw=1.4, ls="--", alpha=0.85)

# percentile guide lines: warm (solid grey) + cold (red) overall p50/p99
for p, ls in ((50, ":"), (99, "-.")):
v = pct(sorted(warm_all), p)
ax.axvline(v, color="#888", ls=ls, lw=0.9)
ax.text(v, 0.04, f"p{p}={v:.1f}ms", rotation=90, fontsize=7,
ax.text(v, 0.04, f"warm p{p}={v:.1f}ms", rotation=90, fontsize=6,
va="bottom", ha="right", color="#555")
if cold_all:
cv = pct(sorted(cold_all), p)
ax.axvline(cv, color="#d62728", ls=ls, lw=0.8, alpha=0.6)
ax.text(cv, 0.04, f"cold p{p}={cv:.1f}ms", rotation=90, fontsize=6,
va="bottom", ha="right", color="#d62728")

ax.set_xlabel("backend query latency (ms)")
ax.set_ylabel("CDF (fraction of queries ≤ x)")
ax.set_title("Fig 7 — backend query latency CDF (warm sketch tier, single-node loopback)")
title = "Fig 7 — backend query latency CDF (single-node loopback)"
if cold_all:
title = ("Fig 7 — backend query latency CDF: warm sketch tier vs "
"cold-fallback archive\n(single-node loopback)")
ax.set_title(title)
ax.set_ylim(0, 1.02)
ax.set_xlim(left=0)
ax.grid(True, alpha=0.3)
ax.legend(loc="lower right", fontsize=8)
ax.legend(loc="lower right", fontsize=7)
fig.tight_layout()
fig.savefig(args.out_png, dpi=140)
print(f"wrote {args.out_json} and {args.out_png}")
print(json.dumps(out["warm"]["overall"], indent=2))
print("WARM:", json.dumps(out["warm"]["overall"]))
for k, v in out["warm"]["by_kind"].items():
print(f" {k:14s} p50={v['p50_ms']:.2f} p95={v['p95_ms']:.2f} p99={v['p99_ms']:.2f} (n={v['n']})")
print(f" warm {k:10s} p50={v['p50_ms']:.2f} p95={v['p95_ms']:.2f} p99={v['p99_ms']:.2f} (n={v['n']})")
if cold_all:
print("COLD:", json.dumps(out["cold"]["overall"]))
for k, v in out["cold"]["by_kind"].items():
print(f" cold {k:10s} p50={v['p50_ms']:.2f} p95={v['p95_ms']:.2f} p99={v['p99_ms']:.2f} (n={v['n']})")


if __name__ == "__main__":
Expand Down
Loading