From ec7a6da045a76e8229b6b0f474d1fd5946cbc6a2 Mon Sep 17 00:00:00 2001 From: Zeying Zhu Date: Fri, 8 May 2026 14:57:40 -0400 Subject: [PATCH] =?UTF-8?q?mvp=5Freducer:=20add=20frequency=20branch=20(ra?= =?UTF-8?q?te(metric[5m]))=20=E2=80=94=20CountMin=20per-series=20=CE=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PR #338 added frequency to the replay client's QUERY_KINDS but the reducer silently skipped frequency rows because parse_query had no `rate()` regex. This adds: - `_RATE_RE` matching `rate(metric[5m])` shape - parse_query returns ("frequency", {"metric": ...}) - new `extract_per_series` helper: PromQL vector → {labels-key → value} - frequency branch in reduce_cell_via_archive: pair warm vs archive per-series, compute mean absolute additive error normalised by truth total — comparable column with rel_err for other kinds - frequency branch in archive_miss fallback path so warm_answer is still captured for diagnostics Refs #46. Closes the gap PR #338 flagged. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.7 (1M context) --- deploy/scripts/accuracy_reduce.py | 67 +++++++++++++++++++++++++++++++ 1 file changed, 67 insertions(+) diff --git a/deploy/scripts/accuracy_reduce.py b/deploy/scripts/accuracy_reduce.py index f263ad03..7ef89512 100755 --- a/deploy/scripts/accuracy_reduce.py +++ b/deploy/scripts/accuracy_reduce.py @@ -331,6 +331,14 @@ def iter_truth_samples(cold_truth_dir: str, metric: str) -> Iterable[dict]: r"sum_over_time\(\s*([\w_]+)\s*\[\s*[0-9smhd]+\s*\]\s*\)", re.IGNORECASE, ) +# frequency shape: `rate(metric[5m])` — per-series rate-of-change. CountMin +# answers this from its frequency estimate (counter increments per window). +# Comparison is against archive's exact rate per series; we report mean +# absolute additive error across keys (the CMS theoretical bound is e/w). +_RATE_RE = re.compile( + r"rate\(\s*([\w_]+)\s*\[\s*[0-9smhd]+\s*\]\s*\)", + re.IGNORECASE, +) def parse_query(promql: str) -> tuple[str, dict] | None: @@ -351,6 +359,8 @@ def parse_query(promql: str) -> tuple[str, dict] | None: return "sum", {"metric": m.group(1)} if (m := _SUM_OVER_TIME_RE.match(s)): return "sum", {"metric": m.group(1)} + if (m := _RATE_RE.match(s)): + return "frequency", {"metric": m.group(1)} return None @@ -430,6 +440,28 @@ def extract_topk_keys(result, k: int) -> list[str]: return keys +def extract_per_series(result) -> dict[str, float]: + """Extract a {labels-key → value} dict from a PromQL vector result. + Used by `frequency` (rate-per-series) comparison: warm CountMin + estimate vs archive exact rate, key-by-key. Returns empty dict on + malformed input.""" + out: dict[str, float] = {} + if not isinstance(result, list): + return out + for el in result: + if not isinstance(el, dict): + continue + labels = el.get("metric") or {} + v = el.get("value") + if not isinstance(v, list) or len(v) < 2: + continue + try: + out[json.dumps(labels, sort_keys=True)] = float(v[1]) + except (TypeError, ValueError): + continue + return out + + # --- per-cell reducer ---------------------------------------------- @@ -556,6 +588,12 @@ def reduce_cell_via_archive( sketch_keys = extract_topk_keys(warm_result, params["k"]) row["warm_answer"] = json.dumps(sketch_keys)[:120] row["answer"] = row["warm_answer"] + elif kind == "frequency": + warm_map = extract_per_series(warm_result) + row["warm_answer"] = json.dumps( + {k: round(v, 4) for k, v in list(warm_map.items())[:5]} + )[:120] + row["answer"] = row["warm_answer"] else: a = extract_scalar(warm_result) if a is not None: @@ -619,6 +657,35 @@ def reduce_cell_via_archive( rel = abs(a - t) / max(abs(t), 1.0) row["rel_err"] = _format_float(rel) row["error"] = row["rel_err"] + elif kind == "frequency": + # CountMin estimates per-key rate; archive returns exact + # per-series rate. Pair by labels-set key, compute mean + # absolute additive error across the intersection. The + # archive_answer / warm_answer columns hold a compact + # JSON snapshot (≤120 char) so the report can show what + # was compared without replaying the query. + truth_map = extract_per_series(archive_result) + warm_map = extract_per_series(warm_result) + row["archive_answer"] = json.dumps( + {k: round(v, 4) for k, v in list(truth_map.items())[:5]} + )[:120] + row["warm_answer"] = json.dumps( + {k: round(v, 4) for k, v in list(warm_map.items())[:5]} + )[:120] + row["truth"] = row["archive_answer"] + row["answer"] = row["warm_answer"] + shared = set(truth_map.keys()) & set(warm_map.keys()) + if shared: + # Mean additive error across overlapping keys, normalised + # by the truth rate's typical magnitude so the column + # stays comparable with rel_err for other kinds. + abs_errs = [abs(warm_map[k] - truth_map[k]) for k in shared] + truth_total = sum(truth_map[k] for k in shared) or 1.0 + rel = (sum(abs_errs) / len(abs_errs)) / max( + truth_total / len(shared), 1.0 + ) + row["rel_err"] = _format_float(rel) + row["error"] = row["rel_err"] writer.writerow(row) n_rows += 1