From b8fa14bcc7751db4c412d2b2c19cad056909c86b Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 17:48:50 -0400 Subject: [PATCH 1/6] refactor(tools): organize post_experiment scripts by scope Split post_experiment/ into single_experiment/ (per-experiment cost, latency, fidelity, throughput analysis), multi_experiment/ (cross- experiment comparison plots), lib/ (results_loader) and debug/. Replace plot_scale_vs_metrics.py with its v2 (adds --experiment_mode) and track plot_latency_cost_tradeoff.py, plot_scale_vs_benefits.py and plot_cardinality_vs_benefit_v2.py. Fix sys.path, imports and sibling script paths for the new layout; update doc references. Co-Authored-By: Claude Opus 5.5 --- asap-tools/README.md | 4 +- .../experiments/post_experiment/README.md | 8 + .../{ => debug}/read_dumped_precomputes.py | 0 .../{ => lib}/results_loader.py | 0 .../README_plot_cardinality_vs_benefit.md | 0 .../plot_cardinality_vs_benefit.py | 15 +- .../plot_cardinality_vs_benefit_v2.py | 709 ++++++++++++++++++ .../plot_comparison_bars.py | 0 .../plot_latency_cost_tradeoff.py | 404 ++++++++++ .../plot_latency_metrics.py | 10 +- .../plot_scale_vs_benefits.py | 386 ++++++++++ .../plot_scale_vs_metrics.py | 113 +-- .../analyze_latencies.py | 9 +- .../analyze_monitor_output.py | 0 .../analyze_throughput.py | 0 .../calculate_fidelity.py | 9 +- .../{ => single_experiment}/compare_costs.py | 4 +- .../compare_latencies.py | 9 +- .../plot_latency_distribution.py | 9 +- .../run_analyze_latencies.sh | 2 +- .../run_calculate_fidelity.sh | 0 .../run_compare_latencies.sh | 2 +- .../run_plot_latency_distribution.sh | 0 .../run-alibaba-cluster-data-experiment.md | 4 +- 24 files changed, 1629 insertions(+), 68 deletions(-) create mode 100644 asap-tools/experiments/post_experiment/README.md rename asap-tools/experiments/post_experiment/{ => debug}/read_dumped_precomputes.py (100%) rename asap-tools/experiments/post_experiment/{ => lib}/results_loader.py (100%) rename asap-tools/experiments/post_experiment/{ => multi_experiment}/README_plot_cardinality_vs_benefit.md (100%) rename asap-tools/experiments/post_experiment/{ => multi_experiment}/plot_cardinality_vs_benefit.py (98%) create mode 100644 asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py rename asap-tools/experiments/post_experiment/{ => multi_experiment}/plot_comparison_bars.py (100%) create mode 100644 asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py rename asap-tools/experiments/post_experiment/{ => multi_experiment}/plot_latency_metrics.py (98%) create mode 100644 asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py rename asap-tools/experiments/post_experiment/{ => multi_experiment}/plot_scale_vs_metrics.py (83%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/analyze_latencies.py (95%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/analyze_monitor_output.py (100%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/analyze_throughput.py (100%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/calculate_fidelity.py (97%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/compare_costs.py (99%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/compare_latencies.py (97%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/plot_latency_distribution.py (97%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/run_analyze_latencies.sh (68%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/run_calculate_fidelity.sh (100%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/run_compare_latencies.sh (91%) rename asap-tools/experiments/post_experiment/{ => single_experiment}/run_plot_latency_distribution.sh (100%) diff --git a/asap-tools/README.md b/asap-tools/README.md index cd0f64c1..d50944e3 100644 --- a/asap-tools/README.md +++ b/asap-tools/README.md @@ -115,7 +115,7 @@ experiments: With this config, 2 experiments are run independently. In the first experiment, `asap-tools/queriers/prometheus-client` only sends queries to ASAP. After this experiment finishes, the infra is torn down. Then the second experiment is set up and `asap-tools/queriers/prometheus-client` sends queries only to Prometheus directly. In the second experiment (i.e. when `mode=prometheus`), none of ASAP's components are set up (apart from `asap-tools/queriers/prometheus-client`). Post-experiment analysis: -- Use `compare_costs.py` and `compare_latencies.py` from `$REPO_DIR/asap-tools/experiments/post_experiments/`. +- Use `compare_costs.py` and `compare_latencies.py` from `$REPO_DIR/asap-tools/experiments/post_experiment/single_experiment/`. - `run_compare_latencies.sh` is an easy wrapper around `compare_latencies.py` ### Comparing query accuracy for ASAP vs Prometheus @@ -130,7 +130,7 @@ experiments: With this config, only one experiment is run. In the same experiment, `PrometheusClient` sends a query to ASAP and then immediately after that, sends a query to Prometheus too. Post-experiment analysis: -- Use `calculate_fidelity.py` from `$REPO_DIR/asap-tools/experiments/post_experiments/`. +- Use `calculate_fidelity.py` from `$REPO_DIR/asap-tools/experiments/post_experiment/single_experiment/`. - `run_calculate_fidelity.sh` is an easy wrapper around `calculate_fidelity.py` ### Debugging with Verbose Logging diff --git a/asap-tools/experiments/post_experiment/README.md b/asap-tools/experiments/post_experiment/README.md new file mode 100644 index 00000000..fb849a35 --- /dev/null +++ b/asap-tools/experiments/post_experiment/README.md @@ -0,0 +1,8 @@ +# post_experiment + +- `single_experiment/` — analyze one experiment (its `baseline` and `sketchdb` modes). Prints stats, `--machine-readable` JSON, optional plots. Cost and latency definitions live here. +- `multi_experiment/` — sweep across experiments and produce comparison figures. Get numbers from `single_experiment/` scripts or `lib/`. +- `lib/` — shared loaders (`results_loader.py`). +- `debug/` — one-off inspection tools. + +Run scripts from any directory; they locate `constants.py` and sibling scripts relative to their own path. diff --git a/asap-tools/experiments/post_experiment/read_dumped_precomputes.py b/asap-tools/experiments/post_experiment/debug/read_dumped_precomputes.py similarity index 100% rename from asap-tools/experiments/post_experiment/read_dumped_precomputes.py rename to asap-tools/experiments/post_experiment/debug/read_dumped_precomputes.py diff --git a/asap-tools/experiments/post_experiment/results_loader.py b/asap-tools/experiments/post_experiment/lib/results_loader.py similarity index 100% rename from asap-tools/experiments/post_experiment/results_loader.py rename to asap-tools/experiments/post_experiment/lib/results_loader.py diff --git a/asap-tools/experiments/post_experiment/README_plot_cardinality_vs_benefit.md b/asap-tools/experiments/post_experiment/multi_experiment/README_plot_cardinality_vs_benefit.md similarity index 100% rename from asap-tools/experiments/post_experiment/README_plot_cardinality_vs_benefit.md rename to asap-tools/experiments/post_experiment/multi_experiment/README_plot_cardinality_vs_benefit.md diff --git a/asap-tools/experiments/post_experiment/plot_cardinality_vs_benefit.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py similarity index 98% rename from asap-tools/experiments/post_experiment/plot_cardinality_vs_benefit.py rename to asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py index 98cb1fff..aa15ed1c 100644 --- a/asap-tools/experiments/post_experiment/plot_cardinality_vs_benefit.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py @@ -45,13 +45,17 @@ ) # Add parent directories to path for imports -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 -from post_experiment.results_loader import ( # noqa: E402 +from post_experiment.lib.results_loader import ( # noqa: E402 load_latencies_only, get_server_name_for_mode, ) -from post_experiment.compare_latencies import calculate_latency_stats # noqa: E402 +from post_experiment.single_experiment.compare_latencies import ( # noqa: E402 + calculate_latency_stats, +) # Metric mapping for cost benefit (compare_costs.py doesn't have 'mean') METRIC_TO_CPU_STAT = { @@ -297,7 +301,10 @@ def extract_cost_benefit_data( ) # Run compare_costs.py - script_dir = os.path.dirname(os.path.abspath(__file__)) + script_dir = os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), + "single_experiment", + ) compare_costs_path = os.path.join(script_dir, "compare_costs.py") result = subprocess.run( diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py new file mode 100644 index 00000000..516ec089 --- /dev/null +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -0,0 +1,709 @@ +#!/usr/bin/env python3 +""" +Plot latency benefit vs lookback period, with one line per cardinality. + +This script analyzes experiments following the naming pattern: + __1_card_2_ + +For example: qot_30m_1_card_2_5 + - query_type: qot (quantile_over_time) + - lookback: 30m + - cardinality: 2^5 = 32 + +The script plots: + - X-axis: Lookback period (log scale, in minutes) + - Y-axis: Latency benefit ratio (prometheus/sketchdb) + - Lines: One per cardinality level (2^0 through 2^9) +""" + +import os +import sys +import re +import glob +import yaml +import argparse +import subprocess +import numpy as np +import pandas as pd +from typing import List, Dict, Any, Optional + +# plotnine imports +from plotnine import ( + ggplot, + aes, + geom_line, + geom_point, + geom_hline, + scale_color_discrete, + scale_x_continuous, + scale_y_continuous, + labs, + theme_minimal, + theme, + element_text, + ggsave, +) + +# Add parent directories to path for imports +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) +import constants # noqa: E402 +from post_experiment.lib.results_loader import ( # noqa: E402 + get_server_name_for_mode, + load_latencies_only, +) +from post_experiment.single_experiment.compare_latencies import ( # noqa: E402 + calculate_latency_stats, +) + +FONTSIZE = 18 + +# Metric mapping for cost benefit (compare_costs.py doesn't have 'mean') +METRIC_TO_CPU_STAT = { + "median": "median", + "p95": "p95", + "p99": "p99", + "sum": "sum", + "max": "max", +} + + +def normalize_lookback(lookback_str: str) -> str: + if lookback_str.endswith("m"): + num_minutes = int(lookback_str[:-1]) + if num_minutes >= 60 and num_minutes % 60 == 0: + hours = num_minutes // 60 + return f"{hours}h" + return lookback_str + + +def parse_lookback_to_minutes(lookback_str: str) -> float: + """ + Convert lookback string to minutes. + + Examples: + '5m' -> 5.0 + '1h' -> 60.0 + '90s' -> 1.5 + '2h30m' -> 150.0 (if needed in future) + """ + lookback_str = lookback_str.strip().lower() + + # Simple patterns first + if lookback_str.endswith("m"): + return float(lookback_str[:-1]) + elif lookback_str.endswith("h"): + return float(lookback_str[:-1]) * 60 + elif lookback_str.endswith("s"): + return float(lookback_str[:-1]) / 60 + + # Fallback: try to parse as just a number (assume minutes) + try: + return float(lookback_str) + except ValueError: + raise ValueError(f"Cannot parse lookback string: {lookback_str}") + + +def parse_experiment_name(exp_name: str) -> Optional[Dict[str, Any]]: + """ + Parse experiment name to extract metadata. + + Expected format: __1_card_2_ + Example: qot_30m_1_card_2_5 + + Returns: + dict with keys: query_type, lookback_str, lookback_minutes, card_exp + or None if name doesn't match pattern + """ + # Pattern: word_lookback_1_card_2_digit + pattern = r"^(?P\w+)_(?P\d+\w+)_1_card_2_(?P\d+)$" + match = re.match(pattern, exp_name) + + if not match: + return None + + data = match.groupdict() + lookback_str = data["lookback"] + + return { + "query_type": data["query_type"], + "lookback_str": lookback_str, + "lookback_minutes": parse_lookback_to_minutes(lookback_str), + "card_exp": int(data["card_exp"]), + } + + +def calculate_data_scale_from_config(config: Dict[str, Any]) -> int: + """ + Calculate data scale from config: num_ports_per_server * num_values_per_label. + """ + fake_exporter_config = config["exporters"]["exporter_list"]["fake_exporter"] + + num_ports_per_server = fake_exporter_config["num_ports_per_server"] + num_values_per_label = fake_exporter_config["num_values_per_label"] + + return num_ports_per_server * num_values_per_label + + +def load_experiment_config(exp_dir: str) -> Dict[str, Any]: + """Load experiment configuration YAML.""" + config_dir = os.path.join(exp_dir, "experiment_config") + + if not os.path.exists(config_dir): + raise FileNotFoundError(f"Config directory not found: {config_dir}") + + config_files = [f for f in os.listdir(config_dir) if f.endswith(".yaml")] + if len(config_files) != 1: + raise ValueError(f"Expected exactly one config file, found {len(config_files)}") + + config_path = os.path.join(config_dir, config_files[0]) + with open(config_path, "r") as f: + return yaml.safe_load(f) + + +def extract_experiment_data( + exp_name: str, metric: str = "p95", verify_scale: bool = True +) -> Optional[Dict[str, Any]]: + """ + Extract data from a single experiment. + + Args: + exp_name: Experiment name + metric: Latency metric to use (median, p95, p99, mean) + verify_scale: If True, verify data scale matches expected 2^card_exp + + Returns: + dict with experiment data or None if extraction fails + """ + # Parse experiment name + metadata = parse_experiment_name(exp_name) + if metadata is None: + print(f"Warning: Skipping {exp_name} (doesn't match naming pattern)") + return None + + exp_dir = os.path.join(constants.LOCAL_EXPERIMENT_DIR, exp_name) + + if not os.path.exists(exp_dir): + print(f"Warning: Experiment directory not found: {exp_dir}") + return None + + try: + # Load config + config = load_experiment_config(exp_dir) + + # Calculate data scale from config + actual_scale = calculate_data_scale_from_config(config) + expected_scale = 2 ** metadata["card_exp"] + + if verify_scale and actual_scale != expected_scale: + print( + f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" + ) + + # Load latencies for both servers + latencies = {} + for server_type in ["baseline", "sketchdb"]: + server_dir = os.path.join(exp_dir, server_type, "prometheus_client_output") + server_name = get_server_name_for_mode(exp_dir, server_type) + + if not os.path.exists(server_dir): + print(f"Warning: {server_type} directory not found for {exp_name}") + return None + + try: + server_latencies = load_latencies_only(server_dir) + if server_name not in server_latencies: + print(f"Warning: No {server_type} data in results for {exp_name}") + return None + + # Aggregate latencies across all queries + all_latencies = [] + for query_idx, latency_result in server_latencies[server_name].items(): + query_latencies = [ + lat for lat in latency_result.get_latencies() if lat is not None + ] + all_latencies.extend(query_latencies) + + if not all_latencies: + print(f"Warning: No latency data for {server_type} in {exp_name}") + return None + + stats = calculate_latency_stats(all_latencies) + latencies[server_type] = stats + + except Exception as e: + print(f"Warning: Failed to load {server_type} data for {exp_name}: {e}") + return None + + # Calculate benefit ratio + if "baseline" not in latencies or "sketchdb" not in latencies: + print(f"Warning: Missing server data for {exp_name}") + return None + + prometheus_latency = latencies["baseline"][metric] + sketchdb_latency = latencies["sketchdb"][metric] + + # Hardcode prometheus median latency for specific experiment + if exp_name == "qot_120m_1_card_2_7" and metric == "median": + prometheus_latency = 0.3 + print( + f"Using hardcoded prometheus median latency: {prometheus_latency} for {exp_name}" + ) + + if sketchdb_latency > 0: + benefit_ratio = prometheus_latency / sketchdb_latency + elif prometheus_latency > 0: + benefit_ratio = float("inf") + else: + benefit_ratio = 1.0 + + return { + "experiment_name": exp_name, + "query_type": metadata["query_type"], + "lookback_str": metadata["lookback_str"], + "lookback_minutes": metadata["lookback_minutes"], + "card_exp": metadata["card_exp"], + "data_scale": actual_scale, + "prometheus_latency": prometheus_latency, + "sketchdb_latency": sketchdb_latency, + "benefit_ratio": benefit_ratio, + "metric": metric, + } + + except Exception as e: + print(f"Warning: Failed to process {exp_name}: {e}") + return None + + +def extract_cost_benefit_data( + exp_name: str, metric: str = "p95", verify_scale: bool = True +) -> Optional[Dict[str, Any]]: + """ + Extract cost benefit data from a single experiment. + + Runs compare_costs.py and parses Query CPU Benefit output. + + Args: + exp_name: Experiment name + metric: CPU metric to use (median, p95, p99, sum, max) + verify_scale: If True, verify data scale matches expected 2^card_exp + + Returns: + dict with experiment data or None if extraction fails + """ + # Parse experiment name + metadata = parse_experiment_name(exp_name) + if metadata is None: + print(f"Warning: Skipping {exp_name} (doesn't match naming pattern)") + return None + + exp_dir = os.path.join(constants.LOCAL_EXPERIMENT_DIR, exp_name) + + if not os.path.exists(exp_dir): + print(f"Warning: Experiment directory not found: {exp_dir}") + return None + + try: + # Verify scale if requested + actual_scale = None + if verify_scale: + config = load_experiment_config(exp_dir) + actual_scale = calculate_data_scale_from_config(config) + expected_scale = 2 ** metadata["card_exp"] + if actual_scale != expected_scale: + print( + f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" + ) + + # Run compare_costs.py + script_dir = os.path.join( + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), + "single_experiment", + ) + compare_costs_path = os.path.join(script_dir, "compare_costs.py") + + result = subprocess.run( + [ + "python3", + compare_costs_path, + "--experiment_name", + exp_name, + "--all_experiment_modes", + "--print", + ], + capture_output=True, + text=True, + check=True, + cwd=script_dir, + ) + + # Parse output for Query CPU Benefit section + output = result.stdout + result.stderr + cpu_stat = METRIC_TO_CPU_STAT.get(metric, "p95") + + # Look for pattern: " : x" in Query CPU Benefit section + # Handle both numeric values and "inf" + pattern = rf"Query CPU Benefit.*?^\s+{cpu_stat}:\s+([\d.]+|inf)x" + match = re.search(pattern, output, re.MULTILINE | re.DOTALL) + + if not match: + print( + f"Warning: Could not find Query CPU Benefit '{cpu_stat}' for {exp_name}" + ) + return None + + value_str = match.group(1) + benefit_ratio = float(value_str) if value_str != "inf" else float("inf") + + return { + "experiment_name": exp_name, + "query_type": metadata["query_type"], + "lookback_str": metadata["lookback_str"], + "lookback_minutes": metadata["lookback_minutes"], + "card_exp": metadata["card_exp"], + "data_scale": actual_scale, + "benefit_ratio": benefit_ratio, + "metric": metric, + } + + except subprocess.CalledProcessError as e: + print(f"Warning: compare_costs.py failed for {exp_name}: {e}") + return None + except Exception as e: + print(f"Warning: Failed to process {exp_name}: {e}") + return None + + +def extract_experiments_from_patterns( + patterns: List[str], + metric: str = "p95", + cardinalities: Optional[List[int]] = None, + benefit_type: str = "latency", + query_types: Optional[List[str]] = None, +) -> pd.DataFrame: + """ + Extract data from experiments matching glob patterns. + + Args: + patterns: List of glob patterns for experiment names + metric: Metric to use (latency or CPU stat depending on benefit_type) + cardinalities: Optional list of cardinality exponents to include + benefit_type: Type of benefit to extract ('latency' or 'cost') + query_types: Optional list of query types to include (e.g., ['qot']) + + Returns: + DataFrame with experiment data + """ + # Find all matching experiment directories + exp_names = set() + for pattern in patterns: + pattern_path = os.path.join(constants.LOCAL_EXPERIMENT_DIR, pattern) + for path in glob.glob(pattern_path): + if os.path.isdir(path): + exp_names.add(os.path.basename(path)) + + if not exp_names: + raise ValueError(f"No experiments found matching patterns: {patterns}") + + print(f"Found {len(exp_names)} experiments matching patterns") + + # Extract data from each experiment - dispatch based on benefit type + data_list = [] + for exp_name in sorted(exp_names): + if benefit_type == "latency": + exp_data = extract_experiment_data(exp_name, metric=metric) + elif benefit_type == "cost": + exp_data = extract_cost_benefit_data(exp_name, metric=metric) + else: + raise ValueError(f"Unknown benefit_type: {benefit_type}") + + if exp_data is not None: + # Filter by cardinality if specified + if cardinalities is not None and exp_data["card_exp"] not in cardinalities: + continue + # Filter by query type if specified + if query_types is not None and exp_data["query_type"] not in query_types: + continue + data_list.append(exp_data) + + if not data_list: + raise ValueError("No valid experiment data extracted") + + print(f"Successfully extracted data from {len(data_list)} experiments") + + df = pd.DataFrame(data_list) + + # Add log-scale transformation: log2(T/15) where T is in minutes + # This makes lookback periods equally spaced on the plot + df["lookback_log2"] = np.log2(df["lookback_minutes"] / 15.0) + + return df + + +def create_plot( + df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" +) -> "ggplot": + """ + Create benefit vs lookback plot with log2(T/15) x-axis. + + Args: + df: DataFrame with experiment data + metric: Metric being plotted (latency or CPU stat) + benefit_type: Type of benefit ('latency' or 'cost') + + Returns: + plotnine ggplot object + """ + # Create labels for legend + card_exps = sorted(df["card_exp"].unique()) + color_labels = [f"2^{exp}" for exp in card_exps] + + # Get unique lookback values for x-axis breaks and labels + lookback_data = df[ + ["lookback_minutes", "lookback_log2", "lookback_str"] + ].drop_duplicates() + lookback_data = lookback_data.sort_values("lookback_minutes") + + x_breaks = lookback_data["lookback_log2"].tolist() + x_labels = lookback_data["lookback_str"].tolist() + + x_labels = [normalize_lookback(lbl) for lbl in x_labels] + + # Get y-axis range and create breaks + y_min = df["benefit_ratio"].min() + y_max = df["benefit_ratio"].max() + + # Create y-axis breaks: use multiples of 10, plus explicitly include 1.0 + y_breaks = [1.0] # Start with 1.0 + step = 10 + current = step + while current <= y_max: + y_breaks.append(float(current)) + current += step + + # Add 0 if needed (if y_min < 1) + if y_min < 1.0: + y_breaks.insert(0, 0.0) + + y_breaks = sorted(list(set(y_breaks))) # Remove duplicates and sort + + # Dynamic Y-axis label based on benefit type + y_label = ( + "Query Latency Benefit" + if benefit_type == "latency" + else "Query Cost Benefit (CPU)" + ) + + p = ( + ggplot( + df, + aes( + x="lookback_log2", + y="benefit_ratio", + color="factor(card_exp)", + group="factor(card_exp)", + ), + ) + + geom_hline(yintercept=1.0, linetype="dashed", color="gray", alpha=0.5) + + geom_line(size=1.2) + + geom_point(size=3) + + scale_x_continuous( + name="Query Lookback Window", breaks=x_breaks, labels=x_labels + ) + + scale_y_continuous(breaks=y_breaks) + + scale_color_discrete(name="Data Cardinality", labels=color_labels) + + labs( + # title=f'{benefit_type.capitalize()} Benefit vs Lookback Period ({metric.upper()})', + y=y_label + ) + + theme_minimal() + + theme( + legend_position="right", + # plot_title=element_text(size=14, weight='bold'), + axis_title_x=element_text(size=FONTSIZE), + axis_title_y=element_text(size=FONTSIZE), + axis_text_x=element_text(size=FONTSIZE), + axis_text_y=element_text(size=FONTSIZE), + legend_title=element_text(size=FONTSIZE), + legend_text=element_text(size=FONTSIZE), + plot_margin=0.1, + ) + ) + + return p + + +def print_summary_table( + df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" +): + """Print summary table of experiment data.""" + # Dynamic header based on benefit type + if benefit_type == "latency": + header = f"Latency Benefit Analysis Summary ({metric.upper()} metric)" + else: + cpu_stat = METRIC_TO_CPU_STAT.get(metric, metric) + header = f"Cost Benefit Analysis Summary (CPU {cpu_stat.upper()})" + + print("\n" + "=" * 100) + print(header) + print("=" * 100) + + # Group by query type + for query_type in df["query_type"].unique(): + query_df = df[df["query_type"] == query_type] + print(f"\nQuery Type: {query_type}") + print("-" * 100) + + # Pivot table: rows = cardinality, columns = lookback + pivot = query_df.pivot_table( + index="card_exp", + columns="lookback_str", + values="benefit_ratio", + aggfunc="mean", + ) + + # Sort columns by lookback minutes + lookback_order = ( + query_df.groupby("lookback_str")["lookback_minutes"].first().sort_values() + ) + pivot = pivot[lookback_order.index] + + # Format with data scale labels + pivot.index = [f"2^{exp} ({2**exp})" for exp in pivot.index] + + print(pivot.to_string(float_format=lambda x: f"{x:.2f}")) + + print("\n" + "=" * 100) + print(f"Total experiments: {len(df)}") + print(f"Cardinalities: {sorted(df['card_exp'].unique())}") + print(f"Lookback periods: {sorted(df['lookback_str'].unique())}") + print(f"Query types: {sorted(df['query_type'].unique())}") + + # Summary statistics + print("\nBenefit Ratio Statistics:") + print(f" Mean: {df['benefit_ratio'].mean():.2f}") + print(f" Median: {df['benefit_ratio'].median():.2f}") + print(f" Min: {df['benefit_ratio'].min():.2f}") + print(f" Max: {df['benefit_ratio'].max():.2f}") + print("=" * 100 + "\n") + + +def main(): + parser = argparse.ArgumentParser( + description="Plot latency or cost benefit vs lookback period, one line per cardinality", + epilog=""" +Examples: + # Print latency benefit summary (default) + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --print + + # Plot and save latency benefit + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --plot --save latency_benefit.png + + # Plot cost benefit with p99 CPU metric + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --benefit-type cost --metric p99 --plot --save cost_benefit_p99.png + + # Plot cost benefit with max CPU metric + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --benefit-type cost --metric max --plot --save cost_benefit_max.png + + # Plot specific cardinalities + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --plot --save benefit.png --cardinalities 0 2 4 6 8 + + # Filter to only 'qot' query type (exclude 'qot_vm' etc.) + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --query-types qot --plot --save benefit.png + + # Print cost benefit summary table + python plot_cardinality_vs_benefit.py "qot_*_1_card_2_*" --benefit-type cost --metric p95 --print + """, + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + + parser.add_argument( + "patterns", + nargs="+", + help='Glob patterns for experiment names (e.g., "qot_*_1_card_2_*")', + ) + parser.add_argument( + "--metric", + type=str, + default="p95", + choices=["median", "p95", "p99", "mean", "sum", "max"], + help="Metric to plot (default: p95). Note: 'mean' only available for latency benefit", + ) + parser.add_argument( + "--benefit-type", + type=str, + default="latency", + choices=["latency", "cost"], + help="Type of benefit to plot: latency or cost (CPU) (default: latency)", + ) + parser.add_argument( + "--cardinalities", + type=int, + nargs="+", + help="Filter to specific cardinality exponents (e.g., 0 2 4 6 8)", + ) + parser.add_argument( + "--query-types", + type=str, + nargs="+", + help="Filter to specific query types (e.g., qot qot_vm)", + ) + parser.add_argument( + "--print", action="store_true", dest="print_summary", help="Print summary table" + ) + parser.add_argument("--plot", action="store_true", help="Generate plot") + parser.add_argument("--save", type=str, help="Save plot to file (provide filename)") + parser.add_argument("--show", action="store_true", help="Display plot") + + args = parser.parse_args() + + # Validate arguments + if args.plot and not (args.save or args.show): + parser.error("--plot requires either --save or --show (or both)") + + if args.save and not args.plot: + parser.error("--save requires --plot") + + if not args.print_summary and not args.plot: + parser.error("Must specify at least one of --print or --plot") + + # Validate metric compatibility with benefit type + if args.benefit_type == "cost" and args.metric == "mean": + parser.error( + "'mean' metric is not available for cost benefit. Use median, p95, p99, sum, or max" + ) + + # Extract experiment data + print( + f"Extracting {args.benefit_type} benefit data from experiments matching: {args.patterns}" + ) + df = extract_experiments_from_patterns( + args.patterns, + metric=args.metric, + cardinalities=args.cardinalities, + benefit_type=args.benefit_type, + query_types=args.query_types, + ) + + # Print summary if requested + if args.print_summary: + print_summary_table(df, metric=args.metric, benefit_type=args.benefit_type) + + # Generate plot if requested + if args.plot: + print("\nGenerating plot...") + plot = create_plot(df, metric=args.metric, benefit_type=args.benefit_type) + + if args.save: + ggsave(plot, args.save, dpi=300, width=10, height=6) + print(f"Plot saved to: {args.save}") + + if args.show: + print(plot) + + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/asap-tools/experiments/post_experiment/plot_comparison_bars.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_comparison_bars.py similarity index 100% rename from asap-tools/experiments/post_experiment/plot_comparison_bars.py rename to asap-tools/experiments/post_experiment/multi_experiment/plot_comparison_bars.py diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py new file mode 100644 index 00000000..5c051116 --- /dev/null +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py @@ -0,0 +1,404 @@ +import os +import sys +import json +import glob +import argparse +import subprocess +import numpy as np +import matplotlib.pyplot as plt +from typing import Dict, Tuple + +POST_EXPERIMENT_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +SINGLE_EXPERIMENT_DIR = os.path.join(POST_EXPERIMENT_DIR, "single_experiment") + +sys.path.append(os.path.dirname(POST_EXPERIMENT_DIR)) +import constants # noqa: E402 + + +def run_compare_costs(experiment_name: str) -> Dict: + """Run compare_costs.py with machine-readable output.""" + script_path = os.path.join(SINGLE_EXPERIMENT_DIR, "compare_costs.py") + + cmd = [ + "python3", + script_path, + "--experiment_name", + experiment_name, + "--all_experiment_modes", + "--print", + "--machine-readable", + ] + + result = subprocess.run(cmd, capture_output=True, text=True) + if result.returncode != 0: + raise RuntimeError( + f"compare_costs.py failed for {experiment_name}: {result.stderr}" + ) + + return json.loads(result.stdout) + + +def run_compare_latencies( + experiment_name: str, exact_mode: str, estimate_mode: str +) -> Dict: + """Run compare_latencies.py with machine-readable output.""" + script_path = os.path.join(SINGLE_EXPERIMENT_DIR, "compare_latencies.py") + + cmd = [ + "python3", + script_path, + "--experiment_name", + experiment_name, + "--exact_experiment_mode", + exact_mode, + "--estimate_experiment_mode", + estimate_mode, + "--machine-readable", + ] + + result = subprocess.run(cmd, capture_output=True, text=True) + if result.returncode != 0: + raise RuntimeError( + f"compare_latencies.py failed for {experiment_name}: {result.stderr}" + ) + + return json.loads(result.stdout) + + +def extract_metrics( + experiment_name: str, + latency_metric: str, + cost_metric: str, + exact_mode: str, + estimate_mode: str, +) -> Tuple[float, float, float, float, float, float]: + """Extract latency and cost metrics for both modes. + + Returns: + (exact_latency, exact_cost, estimate_latency, estimate_cost, + exact_total_cpu, estimate_total_cpu) + """ + # Get cost data + cost_data = run_compare_costs(experiment_name) + + if "query_cpu" not in cost_data: + raise ValueError(f"No query_cpu data found for {experiment_name}") + + if exact_mode not in cost_data["query_cpu"]: + raise ValueError( + f"Mode {exact_mode} not found in query_cpu data for {experiment_name}" + ) + if estimate_mode not in cost_data["query_cpu"]: + raise ValueError( + f"Mode {estimate_mode} not found in query_cpu data for {experiment_name}" + ) + + exact_cost = cost_data["query_cpu"][exact_mode][cost_metric] + estimate_cost = cost_data["query_cpu"][estimate_mode][cost_metric] + + # Total CPU across all monitored processes ("all" pseudo-process in compare_costs.py) + def total_cpu(mode): + return cost_data["experiment_modes"][mode]["processes"]["all_all"][ + "cpu_percent" + ][cost_metric] + + exact_total_cpu = total_cpu(exact_mode) + estimate_total_cpu = total_cpu(estimate_mode) + + # Get latency data + latency_data = run_compare_latencies(experiment_name, exact_mode, estimate_mode) + + if "results" not in latency_data: + raise ValueError(f"No results found for {experiment_name}") + + # Use aggregate results (key "-1" as string since JSON converts int keys to strings) + if "-1" not in latency_data["results"]: + raise ValueError(f"No aggregate results found for {experiment_name}") + + exact_latency = latency_data["results"]["-1"]["exact"][latency_metric] + estimate_latency = latency_data["results"]["-1"]["estimate"][latency_metric] + + return ( + exact_latency, + exact_cost, + estimate_latency, + estimate_cost, + exact_total_cpu, + estimate_total_cpu, + ) + + +def plot_latency_cost_tradeoff( + data_points: Dict[str, Tuple[float, float, float, float]], + latency_metric: str, + cost_metric: str, + args, +): + """Plot latency-cost tradeoff. + + Args: + data_points: Dict mapping experiment_name to (exact_latency, exact_cost, estimate_latency, estimate_cost) + latency_metric: Name of latency metric (e.g., "median") + cost_metric: Name of cost metric (e.g., "p99") + args: Command-line arguments + """ + plt.rcParams.update({"font.size": 24}) + + fig, ax = plt.subplots(figsize=(12, 8)) + + prometheus_latencies = [] + prometheus_costs = [] + turboprom_latencies = [] + turboprom_costs = [] + experiment_names = [] + + for exp_name, (exact_lat, exact_cost, est_lat, est_cost) in data_points.items(): + prometheus_latencies.append(exact_lat) + prometheus_costs.append(exact_cost) + turboprom_latencies.append(est_lat) + turboprom_costs.append(est_cost) + experiment_names.append(exp_name) + + # Plot prometheus points + ax.scatter( + prometheus_latencies, + prometheus_costs, + color="red", + marker="o", + s=100, + alpha=0.6, + label="Prometheus", + ) + + # Plot turboprom points + ax.scatter( + turboprom_latencies, + turboprom_costs, + color="blue", + marker="s", + s=100, + alpha=0.6, + label="ASAPOlly", + ) + + # Optionally label points with experiment names + if args.label_points: + for i, exp_name in enumerate(experiment_names): + # Label prometheus point + ax.annotate( + exp_name, + (prometheus_latencies[i], prometheus_costs[i]), + xytext=(5, 5), + textcoords="offset points", + fontsize=8, + alpha=0.7, + ) + + # Calculate and draw benefit arrows + median_prom_latency = np.median(prometheus_latencies) + median_prom_cost = np.median(prometheus_costs) + median_turbo_latency = np.median(turboprom_latencies) + median_turbo_cost = np.median(turboprom_costs) + + latency_benefit = median_prom_latency / median_turbo_latency + cost_benefit = median_prom_cost / median_turbo_cost + + # Draw horizontal arrow for latency benefit + # Position it at a Y coordinate between the minimum and median cost + min_cost = min(min(prometheus_costs), min(turboprom_costs)) + max_cost = max(max(prometheus_costs), max(turboprom_costs)) + arrow_y_latency = min_cost + (max_cost - min_cost) * 0.15 + ax.annotate( + "", + xy=(median_turbo_latency, arrow_y_latency), + xytext=(median_prom_latency, arrow_y_latency), + arrowprops=dict( + arrowstyle="<->", color="green", lw=2.5, alpha=0.8, shrinkA=0, shrinkB=0 + ), + ) + # Label the latency benefit arrow + ax.text( + (median_prom_latency + median_turbo_latency) / 2, + arrow_y_latency, + f"{latency_benefit:.1f}×", + ha="center", + va="top", + fontsize=24, + fontweight="bold", + color="green", + bbox=dict( + boxstyle="round,pad=0.5", facecolor="white", edgecolor="green", alpha=0.8 + ), + ) + + # Draw vertical arrow for cost benefit + # Position it at an X coordinate between the minimum and median latency + min_latency = min(min(prometheus_latencies), min(turboprom_latencies)) + max_latency = max(max(prometheus_latencies), max(turboprom_latencies)) + arrow_x_cost = min_latency + (max_latency - min_latency) * 0.15 + ax.annotate( + "", + xy=(arrow_x_cost, median_turbo_cost), + xytext=(arrow_x_cost, median_prom_cost), + arrowprops=dict( + arrowstyle="<->", color="purple", lw=2.5, alpha=0.8, shrinkA=0, shrinkB=0 + ), + ) + # Label the cost benefit arrow + ax.text( + arrow_x_cost, + (median_prom_cost + median_turbo_cost) / 2, + f"{cost_benefit:.1f}×", + ha="right", + va="center", + fontsize=24, + fontweight="bold", + color="purple", + bbox=dict( + boxstyle="round,pad=0.5", facecolor="white", edgecolor="purple", alpha=0.8 + ), + ) + + ax.set_xlabel(f"Latency ({latency_metric}) [ms]") + ax.set_ylabel(f"{args.cpu_type.capitalize()} CPU Cost ({cost_metric}) [%]") + ax.set_title("Latency-Cost Tradeoff: Prometheus vs ASAPOlly") + ax.legend() + ax.grid(True, alpha=0.3) + + # Save or show + if args.save: + output_path = args.output_file + plt.savefig(output_path, dpi=300, bbox_inches="tight") + print(f"Saved plot to {output_path}") + + if args.show: + plt.show() + else: + plt.close() + + +def main(args): + if not args.show and not args.save: + raise ValueError("Must specify either --show or --save") + + # Find matching experiment directories + experiment_dirs = glob.glob( + os.path.join(constants.LOCAL_EXPERIMENT_DIR, args.experiment_glob) + ) + + if not experiment_dirs: + raise ValueError(f"No experiments found matching glob: {args.experiment_glob}") + + # Extract experiment names + experiment_names = [os.path.basename(exp_dir) for exp_dir in experiment_dirs] + + print( + f"Found {len(experiment_names)} experiments matching glob: {args.experiment_glob}" + ) + print(f"Experiments: {experiment_names}") + + # Collect data for each experiment + data_points = {} + failed_experiments = [] + + for exp_name in experiment_names: + try: + print(f"\nProcessing experiment: {exp_name}") + ( + exact_lat, + exact_cost, + est_lat, + est_cost, + exact_total, + est_total, + ) = extract_metrics( + exp_name, + args.latency_metric, + args.cost_metric, + args.exact_mode, + args.estimate_mode, + ) + if args.cpu_type == "total": + exact_cost, est_cost = exact_total, est_total + data_points[exp_name] = (exact_lat, exact_cost, est_lat, est_cost) + print( + f" Prometheus: latency={exact_lat:.2f}ms, {args.cpu_type}_cpu={exact_cost:.2f}%" + ) + print( + f" TurboProm: latency={est_lat:.2f}ms, {args.cpu_type}_cpu={est_cost:.2f}%" + ) + except Exception as e: + print(f" Failed to process {exp_name}: {e}") + failed_experiments.append(exp_name) + + if not data_points: + raise ValueError("No valid data points collected") + + if failed_experiments: + print(f"\nWarning: Failed to process {len(failed_experiments)} experiments:") + for exp in failed_experiments: + print(f" - {exp}") + + # Plot the data + plot_latency_cost_tradeoff(data_points, args.latency_metric, args.cost_metric, args) + + +if __name__ == "__main__": + parser = argparse.ArgumentParser( + description="Plot latency-cost tradeoff for multiple experiments" + ) + parser.add_argument( + "--experiment_glob", + type=str, + required=True, + help="Glob pattern to match experiment names (e.g., 'quantile_*')", + ) + parser.add_argument( + "--latency_metric", + type=str, + default="median", + choices=["median", "mean", "p95", "p99"], + help="Latency metric to use (default: median)", + ) + parser.add_argument( + "--cost_metric", + type=str, + default="p99", + choices=["median", "mean", "p95", "p99", "sum"], + help="Cost metric to use (default: p99)", + ) + parser.add_argument( + "--cpu_type", + type=str, + required=True, + choices=["query", "total"], + help="CPU to print/plot: query-attributed CPU or total CPU across all processes", + ) + parser.add_argument( + "--exact_mode", + type=str, + default="baseline", + help="Name of exact/baseline experiment mode (default: baseline)", + ) + parser.add_argument( + "--estimate_mode", + type=str, + default="sketchdb", + help="Name of estimate/optimized experiment mode (default: sketchdb)", + ) + parser.add_argument("--show", action="store_true", help="Show the plot") + parser.add_argument("--save", action="store_true", help="Save the plot to a file") + parser.add_argument( + "--output_file", + type=str, + default="latency_cost_tradeoff.png", + help="Output file path (default: latency_cost_tradeoff.png)", + ) + parser.add_argument( + "--label_points", + action="store_true", + help="Label points with experiment names", + ) + + args = parser.parse_args() + main(args) diff --git a/asap-tools/experiments/post_experiment/plot_latency_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py similarity index 98% rename from asap-tools/experiments/post_experiment/plot_latency_metrics.py rename to asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py index 74090b21..7b1b7c49 100755 --- a/asap-tools/experiments/post_experiment/plot_latency_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py @@ -31,13 +31,17 @@ ) # Add parent directories to path for imports -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 -from post_experiment.results_loader import ( # noqa: E402 +from post_experiment.lib.results_loader import ( # noqa: E402 load_latencies_only, get_server_name_for_mode, ) -from post_experiment.compare_latencies import calculate_latency_stats # noqa: E402 +from post_experiment.single_experiment.compare_latencies import ( # noqa: E402 + calculate_latency_stats, +) class DataExtractor: diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py new file mode 100644 index 00000000..b9ab5cb0 --- /dev/null +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py @@ -0,0 +1,386 @@ +#!/usr/bin/env python3 +""" +Script to plot data scale vs benefits (prometheus/sketchdb ratios). +Shows how much faster and cheaper sketchdb is compared to prometheus. +X-axis: Data scale (metrics/sec) in log scale +Y-axes: Left = Latency Benefit (ratio), Right = Cost Benefit (ratio) +""" + +import argparse +import json +import matplotlib.pyplot as plt +import numpy as np + +# Import functions from the other script +from plot_scale_vs_metrics import ( + calculate_data_scale, + get_latency_p95, + get_cost_p95, + get_query_cost_sum, + get_query_cost_95, +) + +# Configuration +EXPERIMENT_NAMES = [ + "quantile_1s_10queries_10valuesperlabel_2", + "quantile_1s_10queries_20valuesperlabel_2", + "quantile_1s_10queries_30valuesperlabel_2", + "quantile_1s_10queries_40valuesperlabel_2", +] + +FONTSIZE = 24 + + +def print_benefits_summary( + experiments, + data_scales, + latency_benefits, + cost_benefits, + use_query_cost_sum=False, + use_query_cost_95=False, +): + """Print summary of the benefits data.""" + if use_query_cost_sum: + cost_label = "Query Cost Sum Benefit (ratio)" + cost_json_key = "query_cost_sum_benefit_ratio" + elif use_query_cost_95: + cost_label = "Query Cost P95 Benefit (ratio)" + cost_json_key = "query_cost_p95_benefit_ratio" + else: + cost_label = "Cost P95 Benefit (ratio)" + cost_json_key = "cost_p95_benefit_ratio" + + print("\nBenefits Summary (Prometheus / SketchDB):") + print("=" * 110) + print( + f"{'Experiment':<50} {'Data Scale':<20} {'Latency Benefit':<20} {cost_label:<20}" + ) + print("-" * 110) + + for exp, scale, lat_benefit, cost_benefit in zip( + experiments, data_scales, latency_benefits, cost_benefits + ): + scale_str = f"{scale:.2e}" if scale is not None else "N/A" + lat_str = f"{lat_benefit:.2f}x" if lat_benefit is not None else "N/A" + cost_str = f"{cost_benefit:.2f}x" if cost_benefit is not None else "N/A" + print(f"{exp:<50} {scale_str:<20} {lat_str:<20} {cost_str:<20}") + + print("=" * 110) + + # Print json-like structure also + print("\nJSON-like Data Structure:") + data_list = [] + for exp, scale, lat_benefit, cost_benefit in zip( + experiments, data_scales, latency_benefits, cost_benefits + ): + data_list.append( + { + "experiment": exp, + "data_scale_metrics_per_sec": scale, + "latency_benefit_ratio": lat_benefit, + cost_json_key: cost_benefit, + } + ) + + print(json.dumps(data_list, indent=4)) + + +def plot_scale_vs_benefits( + experiments, + data_scales, + latency_benefits, + cost_benefits, + save_file=None, + show=False, + use_query_cost_sum=False, + use_query_cost_95=False, +): + """ + Plot data scale vs benefits (prometheus/sketchdb ratios). + + Args: + experiments: List of experiment names + data_scales: List of data scale values (metrics/sec) + latency_benefits: List of latency benefit ratios (prometheus/sketchdb) + cost_benefits: List of cost benefit ratios (prometheus/sketchdb) + save_file: Filename to save the plot (if None, doesn't save) + show: Whether to display the plot + use_query_cost_sum: Whether cost values represent query cost sum + use_query_cost_95: Whether cost values represent query cost p95 + + Returns: + matplotlib figure object + """ + # Filter out None values and sort by data scale + valid_data = [ + (s, l, c, e) + for s, l, c, e in zip(data_scales, latency_benefits, cost_benefits, experiments) + if s is not None and l is not None and c is not None + ] + + if not valid_data: + print("Error: No valid data points to plot") + return None + + valid_data.sort(key=lambda x: x[0]) # Sort by data scale + data_scales_sorted, latency_benefits_sorted, cost_benefits_sorted, _ = zip( + *valid_data + ) + + # Convert to numpy arrays + data_scales_arr = np.array(data_scales_sorted) + latency_benefits_arr = np.array(latency_benefits_sorted) + cost_benefits_arr = np.array(cost_benefits_sorted) + + # Create the plot with two y-axes + fig, ax1 = plt.subplots(figsize=(12, 6)) + + # Determine cost label based on type + if use_query_cost_sum: + cost_ylabel = "Query Cost Benefit (ratio)" + cost_legend = "Query Cost Benefit" + elif use_query_cost_95: + cost_ylabel = "Query Cost P95 Benefit (ratio)" + cost_legend = "Query Cost P95 Benefit" + else: + cost_ylabel = "Cost P95 Benefit (ratio)" + cost_legend = "Cost P95 Benefit" + + # Plot cost benefit on left y-axis + color_cost = "#1f77b4" + # ax1.set_xlabel("Data Scale (metrics/sec)", fontsize=FONTSIZE, fontweight="bold") + ax1.set_xlabel("Data Cardinality", fontsize=FONTSIZE, fontweight="bold") + ax1.set_ylabel(cost_ylabel, fontsize=FONTSIZE, fontweight="bold", color=color_cost) + line1 = ax1.plot( + data_scales_arr, + cost_benefits_arr, + "o-", + color=color_cost, + linewidth=2, + markersize=8, + label=cost_legend, + ) + ax1.tick_params(axis="y", labelcolor=color_cost, labelsize=FONTSIZE) + ax1.tick_params(axis="x", labelsize=FONTSIZE) + ax1.set_xscale("log") + ax1.grid(True, alpha=0.3, which="both") + + # Create second y-axis for latency benefit + ax2 = ax1.twinx() + color_latency = "#ff7f0e" + ax2.set_ylabel( + "Latency Benefit (ratio)", + fontsize=FONTSIZE, + fontweight="bold", + color=color_latency, + ) + line2 = ax2.plot( + data_scales_arr, + latency_benefits_arr, + "s-", + color=color_latency, + linewidth=2, + markersize=8, + label="Latency Benefit", + ) + ax2.tick_params(axis="y", labelcolor=color_latency, labelsize=FONTSIZE) + + # Add title + plt.title( + "TurboProm's Benefits vs Data Cardinality", + fontsize=FONTSIZE + 2, + fontweight="bold", + pad=20, + ) + + # Add vertical dotted lines and annotations for each data point + # Calculate the middle position for annotations (in data coordinates) + y1_min, y1_max = ax1.get_ylim() + annotation_y = y1_min + (y1_max - y1_min) * 0.5 # Center vertically + + for i, x in enumerate(data_scales_arr): + # Format the data scale nicely + if x < 1000: + scale_label = f"{int(x)}" + elif x < 1000000: + scale_label = f"{int(x/1000)}K" + else: + scale_label = f"{x/1000000:.1f}M" + + # Draw vertical dotted line + ax1.axvline(x=x, color="gray", linestyle=":", alpha=0.5, linewidth=1.5) + + # Add annotation at the center of the plot + ax1.text( + x, + annotation_y, + scale_label, + ha="center", + va="center", + fontsize=FONTSIZE - 4, + bbox=dict( + boxstyle="round,pad=0.4", facecolor="white", edgecolor="gray", alpha=0.8 + ), + ) + + # Add legend + lines = line1 + line2 + labels = [line.get_label() for line in lines] + ax1.legend(lines, labels, loc="upper left", fontsize=FONTSIZE) + + # Adjust layout + fig.tight_layout() + + # Save if requested + if save_file: + plt.savefig(save_file, dpi=300, bbox_inches="tight") + print(f"Plot saved as '{save_file}'") + + # Show if requested + if show: + plt.show() + else: + plt.close(fig) + + return fig + + +def main(): + parser = argparse.ArgumentParser( + description="Plot data scale vs benefits (prometheus/sketchdb ratios)", + epilog=""" +Examples: + # Print benefits summary only + python3 plot_scale_vs_benefits.py --print + + # Plot and save to file + python3 plot_scale_vs_benefits.py --plot --save scale_benefits.png + + # Plot and show interactively + python3 plot_scale_vs_benefits.py --plot --show + + # Both print and plot + python3 plot_scale_vs_benefits.py --print --plot --save output.png --show + + # Use query cost sum instead of p95 cost + python3 plot_scale_vs_benefits.py --print --use-query-cost-sum + """, + formatter_class=argparse.RawDescriptionHelpFormatter, + ) + + parser.add_argument("--print", action="store_true", help="Print benefits summary") + parser.add_argument("--plot", action="store_true", help="Generate plot") + parser.add_argument( + "--save", + type=str, + metavar="FILENAME", + help="Save plot to file (provide filename)", + ) + parser.add_argument("--show", action="store_true", help="Display plot") + parser.add_argument( + "--use-query-cost-sum", + action="store_true", + help="Use query CPU cost sum instead of p95 CPU cost", + ) + parser.add_argument( + "--use-query-cost-95", + action="store_true", + help="Use query CPU cost p95 instead of p95 CPU cost", + ) + + args = parser.parse_args() + + # Validate arguments + if args.plot and not (args.save or args.show): + parser.error("--plot requires either --save or --show (or both)") + + if not args.print and not args.plot: + parser.error("At least one of --print or --plot must be specified") + + # Collect data for all experiments + print(f"Processing {len(EXPERIMENT_NAMES)} experiments...") + + data_scales = [] + latency_benefits = [] + cost_benefits = [] + + for exp_name in EXPERIMENT_NAMES: + print(f"\nProcessing: {exp_name}") + + # Calculate data scale + scale = calculate_data_scale(exp_name) + data_scales.append(scale) + if scale is not None: + print(f" Data scale: {scale:.2e} metrics/sec") + + # Get latency for both prometheus and sketchdb + latency_prometheus = get_latency_p95(exp_name, "baseline") + latency_sketchdb = get_latency_p95(exp_name, "sketchdb") + + if latency_prometheus is not None and latency_sketchdb is not None: + latency_benefit = latency_prometheus / latency_sketchdb + latency_benefits.append(latency_benefit) + print( + f" Latency benefit: {latency_benefit:.2f}x (prometheus: {latency_prometheus:.4f}s, sketchdb: {latency_sketchdb:.4f}s)" + ) + else: + latency_benefits.append(None) + print( + f" Latency benefit: N/A (prometheus: {latency_prometheus}, sketchdb: {latency_sketchdb})" + ) + + # Get cost for both prometheus and sketchdb + if args.use_query_cost_sum: + cost_prometheus = get_query_cost_sum(exp_name, "baseline") + cost_sketchdb = get_query_cost_sum(exp_name, "sketchdb") + cost_type = "Query cost sum" + elif args.use_query_cost_95: + cost_prometheus = get_query_cost_95(exp_name, "baseline") + cost_sketchdb = get_query_cost_95(exp_name, "sketchdb") + cost_type = "Query cost p95" + else: + cost_prometheus = get_cost_p95(exp_name, "baseline") + cost_sketchdb = get_cost_p95(exp_name, "sketchdb") + cost_type = "Cost p95" + + if cost_prometheus is not None and cost_sketchdb is not None: + cost_benefit = cost_prometheus / cost_sketchdb + cost_benefits.append(cost_benefit) + print( + f" {cost_type} benefit: {cost_benefit:.2f}x (prometheus: {cost_prometheus:.2f}%, sketchdb: {cost_sketchdb:.2f}%)" + ) + else: + cost_benefits.append(None) + print( + f" {cost_type} benefit: N/A (prometheus: {cost_prometheus}, sketchdb: {cost_sketchdb})" + ) + + # Print summary if requested + if args.print: + print_benefits_summary( + EXPERIMENT_NAMES, + data_scales, + latency_benefits, + cost_benefits, + use_query_cost_sum=args.use_query_cost_sum, + use_query_cost_95=args.use_query_cost_95, + ) + + # Generate plot if requested + if args.plot: + plot_scale_vs_benefits( + experiments=EXPERIMENT_NAMES, + data_scales=data_scales, + latency_benefits=latency_benefits, + cost_benefits=cost_benefits, + save_file=args.save, + show=args.show, + use_query_cost_sum=args.use_query_cost_sum, + use_query_cost_95=args.use_query_cost_95, + ) + + return 0 + + +if __name__ == "__main__": + exit(main()) diff --git a/asap-tools/experiments/post_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py similarity index 83% rename from asap-tools/experiments/post_experiment/plot_scale_vs_metrics.py rename to asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index fb4c72e8..5526ee04 100755 --- a/asap-tools/experiments/post_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -15,7 +15,10 @@ import matplotlib.pyplot as plt import numpy as np -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +POST_EXPERIMENT_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +SINGLE_EXPERIMENT_DIR = os.path.join(POST_EXPERIMENT_DIR, "single_experiment") + +sys.path.append(os.path.dirname(POST_EXPERIMENT_DIR)) import constants # noqa: E402 # Configuration @@ -83,17 +86,18 @@ def calculate_data_scale(experiment_name): return None -def get_latency_p95(experiment_name): +def get_latency_p95(experiment_name, experiment_mode="baseline"): """ Get p95 latency by running ./run_compare_latencies.sh Args: experiment_name: Name of the experiment + experiment_mode: Experiment mode (baseline or sketchdb) Returns: p95 latency value (exact), or None if failed """ - script_dir = os.path.dirname(os.path.abspath(__file__)) + script_dir = SINGLE_EXPERIMENT_DIR script_path = os.path.join(script_dir, "run_compare_latencies.sh") try: @@ -105,20 +109,28 @@ def get_latency_p95(experiment_name): cwd=script_dir, ) - # Parse output to extract p95 from exact - # Looking for: exact: {'median': X, 'p95': Y, ...} + # Parse output to extract p95 from the appropriate line based on experiment mode + # For baseline: look for "exact:" line + # For sketchdb: look for "estimate:" line output = result.stdout + result.stderr - # Find the "exact:" line - exact_match = re.search(r"exact:\s*\{([^}]+)\}", output) - if exact_match: - exact_dict_str = exact_match.group(1) + # Map experiment mode to the line to search for + search_key = "exact" if experiment_mode == "baseline" else "estimate" + + # Find the line with the appropriate key + pattern = rf"{search_key}:\s*\{{([^}}]+)\}}" + match = re.search(pattern, output) + + if match: + dict_str = match.group(1) # Extract p95 value - p95_match = re.search(r"'p95':\s*([\d.]+)", exact_dict_str) + p95_match = re.search(r"'p95':\s*([\d.]+)", dict_str) if p95_match: return float(p95_match.group(1)) - print(f"Warning: Could not parse p95 latency from output for {experiment_name}") + print( + f"Warning: Could not parse p95 latency from output for {experiment_name} in {experiment_mode} mode" + ) return None except subprocess.CalledProcessError as e: @@ -129,17 +141,18 @@ def get_latency_p95(experiment_name): return None -def get_cost_p95(experiment_name): +def get_cost_p95(experiment_name, experiment_mode="baseline"): """ Get p95 CPU cost by running compare_costs.py Args: experiment_name: Name of the experiment + experiment_mode: Experiment mode (baseline or sketchdb) Returns: p95 CPU percentage, or None if failed """ - script_dir = os.path.dirname(os.path.abspath(__file__)) + script_dir = SINGLE_EXPERIMENT_DIR compare_costs_path = os.path.join(script_dir, "compare_costs.py") try: @@ -150,7 +163,7 @@ def get_cost_p95(experiment_name): "--experiment_name", experiment_name, "--experiment_mode", - "baseline", + experiment_mode, "--print", ], capture_output=True, @@ -159,18 +172,17 @@ def get_cost_p95(experiment_name): cwd=script_dir, ) - # Parse output to extract p95 CPU from "prometheus prometheus.yml cpu_percent p95" - # or "prometheus prometheus cpu_percent p95" + # Parse output to extract p95 CPU from "{experiment_mode} {experiment_mode}.yml cpu_percent p95" + # or "{experiment_mode} {experiment_mode} cpu_percent p95" output = result.stdout + result.stderr # Look for lines matching the pattern + pattern = ( + rf"{experiment_mode}\s+{experiment_mode}.*cpu_percent\s+p95\s+([\d.]+)" + ) for line in output.split("\n"): - if re.search( - r"prometheus\s+prometheus.*cpu_percent\s+p95\s+([\d.]+)", line - ): - match = re.search( - r"prometheus\s+prometheus.*cpu_percent\s+p95\s+([\d.]+)", line - ) + if re.search(pattern, line): + match = re.search(pattern, line) if match: return float(match.group(1)) @@ -187,17 +199,18 @@ def get_cost_p95(experiment_name): return None -def get_query_cost_95(experiment_name): +def get_query_cost_95(experiment_name, experiment_mode="baseline"): """ Get query CPU cost p95 by running compare_costs.py Args: experiment_name: Name of the experiment + experiment_mode: Experiment mode (baseline or sketchdb) Returns: Query CPU cost p95 percentage, or None if failed """ - script_dir = os.path.dirname(os.path.abspath(__file__)) + script_dir = SINGLE_EXPERIMENT_DIR compare_costs_path = os.path.join(script_dir, "compare_costs.py") try: @@ -208,7 +221,7 @@ def get_query_cost_95(experiment_name): "--experiment_name", experiment_name, "--experiment_mode", - "baseline", + experiment_mode, "--print", ], capture_output=True, @@ -219,26 +232,26 @@ def get_query_cost_95(experiment_name): # Parse output to extract query CPU sum from "Query CPU Statistics" section # Looking for pattern like: - # prometheus: + # {experiment_mode}: # p95: 1122837.55% output = result.stdout + result.stderr # Look for the Query CPU Statistics section in_query_section = False - in_prometheus_subsection = False + in_experiment_mode_subsection = False for line in output.split("\n"): if "Query CPU Statistics" in line: in_query_section = True continue if in_query_section: - # Check if we're in the prometheus subsection - if line.strip().startswith("prometheus:"): - in_prometheus_subsection = True + # Check if we're in the experiment_mode subsection + if line.strip().startswith(f"{experiment_mode}:"): + in_experiment_mode_subsection = True continue - # If we're in prometheus subsection, look for sum - if in_prometheus_subsection: + # If we're in experiment_mode subsection, look for sum + if in_experiment_mode_subsection: match = re.search(r"p95:\s+([\d.]+)%", line) if match: return float(match.group(1)) @@ -261,17 +274,18 @@ def get_query_cost_95(experiment_name): return None -def get_query_cost_sum(experiment_name): +def get_query_cost_sum(experiment_name, experiment_mode="baseline"): """ Get query CPU cost sum by running compare_costs.py Args: experiment_name: Name of the experiment + experiment_mode: Experiment mode (baseline or sketchdb) Returns: Query CPU cost sum percentage, or None if failed """ - script_dir = os.path.dirname(os.path.abspath(__file__)) + script_dir = SINGLE_EXPERIMENT_DIR compare_costs_path = os.path.join(script_dir, "compare_costs.py") try: @@ -282,7 +296,7 @@ def get_query_cost_sum(experiment_name): "--experiment_name", experiment_name, "--experiment_mode", - "baseline", + experiment_mode, "--print", ], capture_output=True, @@ -293,26 +307,26 @@ def get_query_cost_sum(experiment_name): # Parse output to extract query CPU sum from "Query CPU Statistics" section # Looking for pattern like: - # prometheus: + # {experiment_mode}: # sum: 1122837.55% output = result.stdout + result.stderr # Look for the Query CPU Statistics section in_query_section = False - in_prometheus_subsection = False + in_experiment_mode_subsection = False for line in output.split("\n"): if "Query CPU Statistics" in line: in_query_section = True continue if in_query_section: - # Check if we're in the prometheus subsection - if line.strip().startswith("prometheus:"): - in_prometheus_subsection = True + # Check if we're in the experiment_mode subsection + if line.strip().startswith(f"{experiment_mode}:"): + in_experiment_mode_subsection = True continue - # If we're in prometheus subsection, look for sum - if in_prometheus_subsection: + # If we're in experiment_mode subsection, look for sum + if in_experiment_mode_subsection: match = re.search(r"sum:\s+([\d.]+)%", line) if match: return float(match.group(1)) @@ -539,6 +553,13 @@ def main(): action="store_true", help="Use query CPU cost p95 instead of p95 CPU cost", ) + parser.add_argument( + "--experiment_mode", + type=str, + choices=["baseline", "sketchdb"], + default="baseline", + help="Experiment mode (baseline or sketchdb)", + ) args = parser.parse_args() @@ -566,22 +587,22 @@ def main(): print(f" Data scale: {scale:.2e} metrics/sec") # Get latency p95 - latency = get_latency_p95(exp_name) + latency = get_latency_p95(exp_name, args.experiment_mode) latencies.append(latency) if latency is not None: print(f" Latency p95: {latency:.4f} seconds") # Get cost (either p95 or query cost sum based on flag) if args.use_query_cost_sum: - cost = get_query_cost_sum(exp_name) + cost = get_query_cost_sum(exp_name, args.experiment_mode) if cost is not None: print(f" Query cost sum: {cost:.2f} CPU %") elif args.use_query_cost_95: - cost = get_query_cost_95(exp_name) + cost = get_query_cost_95(exp_name, args.experiment_mode) if cost is not None: print(f" Query cost p95: {cost:.2f} CPU %") else: - cost = get_cost_p95(exp_name) + cost = get_cost_p95(exp_name, args.experiment_mode) if cost is not None: print(f" Cost p95: {cost:.2f} CPU %") costs.append(cost) diff --git a/asap-tools/experiments/post_experiment/analyze_latencies.py b/asap-tools/experiments/post_experiment/single_experiment/analyze_latencies.py similarity index 95% rename from asap-tools/experiments/post_experiment/analyze_latencies.py rename to asap-tools/experiments/post_experiment/single_experiment/analyze_latencies.py index 5f6f8d06..3c1e32cf 100644 --- a/asap-tools/experiments/post_experiment/analyze_latencies.py +++ b/asap-tools/experiments/post_experiment/single_experiment/analyze_latencies.py @@ -9,7 +9,9 @@ from promql_utilities.query_results.classes import LatencyResultAcrossTime # TODO: make this more robust -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 @@ -106,7 +108,10 @@ def print_analysis_results( def main(args): experiment_dir = os.path.join(constants.LOCAL_EXPERIMENT_DIR, args.experiment_name) - from results_loader import load_latencies_only, get_server_name_for_mode + from post_experiment.lib.results_loader import ( + load_latencies_only, + get_server_name_for_mode, + ) if not args.experiment_server_name: args.experiment_server_name = get_server_name_for_mode( diff --git a/asap-tools/experiments/post_experiment/analyze_monitor_output.py b/asap-tools/experiments/post_experiment/single_experiment/analyze_monitor_output.py similarity index 100% rename from asap-tools/experiments/post_experiment/analyze_monitor_output.py rename to asap-tools/experiments/post_experiment/single_experiment/analyze_monitor_output.py diff --git a/asap-tools/experiments/post_experiment/analyze_throughput.py b/asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py similarity index 100% rename from asap-tools/experiments/post_experiment/analyze_throughput.py rename to asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py diff --git a/asap-tools/experiments/post_experiment/calculate_fidelity.py b/asap-tools/experiments/post_experiment/single_experiment/calculate_fidelity.py similarity index 97% rename from asap-tools/experiments/post_experiment/calculate_fidelity.py rename to asap-tools/experiments/post_experiment/single_experiment/calculate_fidelity.py index 6c773a84..80525ebf 100644 --- a/asap-tools/experiments/post_experiment/calculate_fidelity.py +++ b/asap-tools/experiments/post_experiment/single_experiment/calculate_fidelity.py @@ -9,7 +9,9 @@ # from promql_utilities.query_results.classes import QueryResult, QueryResultAcrossTime # TODO: make this more robust -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 @@ -173,7 +175,10 @@ def main(args): exact_results = None estimate_results = None - from results_loader import load_results, get_server_name_for_mode + from post_experiment.lib.results_loader import ( + load_results, + get_server_name_for_mode, + ) exact_results = load_results( os.path.join( diff --git a/asap-tools/experiments/post_experiment/compare_costs.py b/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py similarity index 99% rename from asap-tools/experiments/post_experiment/compare_costs.py rename to asap-tools/experiments/post_experiment/single_experiment/compare_costs.py index 1984f2d3..b5f1c68f 100644 --- a/asap-tools/experiments/post_experiment/compare_costs.py +++ b/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py @@ -8,7 +8,9 @@ from typing import List from collections import defaultdict -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 RESOURCES = ["cpu_percent", "memory_info"] diff --git a/asap-tools/experiments/post_experiment/compare_latencies.py b/asap-tools/experiments/post_experiment/single_experiment/compare_latencies.py similarity index 97% rename from asap-tools/experiments/post_experiment/compare_latencies.py rename to asap-tools/experiments/post_experiment/single_experiment/compare_latencies.py index 8de9a2b6..42c01250 100644 --- a/asap-tools/experiments/post_experiment/compare_latencies.py +++ b/asap-tools/experiments/post_experiment/single_experiment/compare_latencies.py @@ -10,7 +10,9 @@ from promql_utilities.query_results.classes import LatencyResultAcrossTime # TODO: make this more robust -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 @@ -156,7 +158,10 @@ def main(args): exact_results: Optional[Dict[int, LatencyResultAcrossTime]] = None estimate_results: Optional[Dict[int, LatencyResultAcrossTime]] = None - from results_loader import load_latencies_only, get_server_name_for_mode + from post_experiment.lib.results_loader import ( + load_latencies_only, + get_server_name_for_mode, + ) import logging # Suppress debug logging in machine-readable mode diff --git a/asap-tools/experiments/post_experiment/plot_latency_distribution.py b/asap-tools/experiments/post_experiment/single_experiment/plot_latency_distribution.py similarity index 97% rename from asap-tools/experiments/post_experiment/plot_latency_distribution.py rename to asap-tools/experiments/post_experiment/single_experiment/plot_latency_distribution.py index d2d75784..aa80ad32 100644 --- a/asap-tools/experiments/post_experiment/plot_latency_distribution.py +++ b/asap-tools/experiments/post_experiment/single_experiment/plot_latency_distribution.py @@ -11,7 +11,9 @@ from promql_utilities.query_results.classes import LatencyResultAcrossTime # TODO: make this more robust -sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.append( + os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +) import constants # noqa: E402 @@ -128,7 +130,10 @@ def print_percentile_data( def main(args): experiment_dir = os.path.join(constants.LOCAL_EXPERIMENT_DIR, args.experiment_name) - from results_loader import load_latencies_only, get_server_name_for_mode + from post_experiment.lib.results_loader import ( + load_latencies_only, + get_server_name_for_mode, + ) if not args.exact_experiment_server_name: args.exact_experiment_server_name = get_server_name_for_mode( diff --git a/asap-tools/experiments/post_experiment/run_analyze_latencies.sh b/asap-tools/experiments/post_experiment/single_experiment/run_analyze_latencies.sh similarity index 68% rename from asap-tools/experiments/post_experiment/run_analyze_latencies.sh rename to asap-tools/experiments/post_experiment/single_experiment/run_analyze_latencies.sh index 24dc2bd8..4e27dbc0 100755 --- a/asap-tools/experiments/post_experiment/run_analyze_latencies.sh +++ b/asap-tools/experiments/post_experiment/single_experiment/run_analyze_latencies.sh @@ -12,4 +12,4 @@ EXP_NAME=$1 EXP_MODE=$2 PER_QUERY_FLAG=$3 -python3 $THIS_DIR/analyze_latencies.py --experiment_name $EXP_NAME --experiment_mode $EXP_MODE ${PER_QUERY_FLAG} +python3 "$THIS_DIR"/analyze_latencies.py --experiment_name "$EXP_NAME" --experiment_mode "$EXP_MODE" ${PER_QUERY_FLAG:+"$PER_QUERY_FLAG"} diff --git a/asap-tools/experiments/post_experiment/run_calculate_fidelity.sh b/asap-tools/experiments/post_experiment/single_experiment/run_calculate_fidelity.sh similarity index 100% rename from asap-tools/experiments/post_experiment/run_calculate_fidelity.sh rename to asap-tools/experiments/post_experiment/single_experiment/run_calculate_fidelity.sh diff --git a/asap-tools/experiments/post_experiment/run_compare_latencies.sh b/asap-tools/experiments/post_experiment/single_experiment/run_compare_latencies.sh similarity index 91% rename from asap-tools/experiments/post_experiment/run_compare_latencies.sh rename to asap-tools/experiments/post_experiment/single_experiment/run_compare_latencies.sh index a6867b09..b1b5abf7 100755 --- a/asap-tools/experiments/post_experiment/run_compare_latencies.sh +++ b/asap-tools/experiments/post_experiment/single_experiment/run_compare_latencies.sh @@ -11,4 +11,4 @@ EXP_NAME=$1 PER_QUERY_FLAG=$2 #python3 compare_latencies.py --experiment_name $EXP_NAME --exact_experiment_mode sketchdb --exact_experiment_server_name baseline --estimate_experiment_mode sketchdb ${PER_QUERY_FLAG} -python3 "$THIS_DIR"/compare_latencies.py --experiment_name "$EXP_NAME" --exact_experiment_mode baseline --estimate_experiment_mode sketchdb ${PER_QUERY_FLAG} +python3 "$THIS_DIR"/compare_latencies.py --experiment_name "$EXP_NAME" --exact_experiment_mode baseline --estimate_experiment_mode sketchdb ${PER_QUERY_FLAG:+"$PER_QUERY_FLAG"} diff --git a/asap-tools/experiments/post_experiment/run_plot_latency_distribution.sh b/asap-tools/experiments/post_experiment/single_experiment/run_plot_latency_distribution.sh similarity index 100% rename from asap-tools/experiments/post_experiment/run_plot_latency_distribution.sh rename to asap-tools/experiments/post_experiment/single_experiment/run_plot_latency_distribution.sh diff --git a/docs/03-how-to-guides/operations/run-alibaba-cluster-data-experiment.md b/docs/03-how-to-guides/operations/run-alibaba-cluster-data-experiment.md index 783f4dbe..582de2e3 100644 --- a/docs/03-how-to-guides/operations/run-alibaba-cluster-data-experiment.md +++ b/docs/03-how-to-guides/operations/run-alibaba-cluster-data-experiment.md @@ -190,7 +190,7 @@ Because Google has separate baseline and sketchdb modes, the latency wrapper can be used after the run: ```bash -cd /home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-tools/experiments/post_experiment +cd /home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-tools/experiments/post_experiment/single_experiment PYTHONPATH=/home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-common/dependencies/py/promql_utilities \ ./run_compare_latencies.sh google_cluster_data_offset18 @@ -294,7 +294,7 @@ modes. The Alibaba example has one `sketchdb` mode with both Prometheus and SketchDB servers, so compare those servers directly: ```bash -cd /home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-tools/experiments/post_experiment +cd /home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-tools/experiments/post_experiment/single_experiment PYTHONPATH=/home/milind/Desktop/cmu/research/sketch_db_for_prometheus/code/ASAPQuery/asap-common/dependencies/py/promql_utilities \ python3 compare_latencies.py \ From 7d52997a60d40659a35afec5af9d71799467bd68 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 17:57:39 -0400 Subject: [PATCH 2/6] refactor(tools): read post_experiment costs from JSON, label total vs query CPU Replace regex parsing of compare_costs.py / run_compare_latencies.sh text output in the multi_experiment scripts with --machine-readable JSON. plot_scale_vs_metrics' default cost is now total CPU p95 (it previously never matched for baseline and matched only the query engine for sketchdb). Add --benefit-type total_cost to the cardinality plots and name total vs query CPU in all cost labels. Co-Authored-By: Claude Opus 5.5 --- .../plot_cardinality_vs_benefit.py | 53 +-- .../plot_cardinality_vs_benefit_v2.py | 59 +-- .../plot_scale_vs_benefits.py | 22 +- .../multi_experiment/plot_scale_vs_metrics.py | 345 ++++-------------- 4 files changed, 145 insertions(+), 334 deletions(-) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py index aa15ed1c..b8020e46 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py @@ -19,6 +19,7 @@ import os import sys import re +import json import glob import yaml import argparse @@ -261,17 +262,18 @@ def extract_experiment_data( def extract_cost_benefit_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True + exp_name: str, metric: str = "p95", verify_scale: bool = True, total: bool = False ) -> Optional[Dict[str, Any]]: """ Extract cost benefit data from a single experiment. - Runs compare_costs.py and parses Query CPU Benefit output. + Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. Args: exp_name: Experiment name metric: CPU metric to use (median, p95, p99, sum, max) verify_scale: If True, verify data scale matches expected 2^card_exp + total: Use total CPU (all processes) instead of query CPU Returns: dict with experiment data or None if extraction fails @@ -315,6 +317,7 @@ def extract_cost_benefit_data( exp_name, "--all_experiment_modes", "--print", + "--machine-readable", ], capture_output=True, text=True, @@ -322,22 +325,16 @@ def extract_cost_benefit_data( cwd=script_dir, ) - # Parse output for Query CPU Benefit section - output = result.stdout + result.stderr + costs = json.loads(result.stdout) cpu_stat = METRIC_TO_CPU_STAT.get(metric, "p95") - - # Look for pattern: " : x" in Query CPU Benefit section - pattern = rf"Query CPU Benefit.*?^\s+{cpu_stat}:\s+([\d.]+)x" - match = re.search(pattern, output, re.MULTILINE | re.DOTALL) - - if not match: - print( - f"Warning: Could not find Query CPU Benefit '{cpu_stat}' for {exp_name}" - ) + benefits = ( + costs["benefit"]["cpu_percent"] if total else costs["query_cpu_benefit"] + ) + benefit_ratio = benefits[cpu_stat] + if np.isnan(benefit_ratio): # 0/0: both modes had zero cost + print(f"Warning: {cpu_stat} CPU benefit is NaN for {exp_name}") return None - benefit_ratio = float(match.group(1)) - return { "experiment_name": exp_name, "query_type": metadata["query_type"], @@ -370,7 +367,7 @@ def extract_experiments_from_patterns( patterns: List of glob patterns for experiment names metric: Metric to use (latency or CPU stat depending on benefit_type) cardinalities: Optional list of cardinality exponents to include - benefit_type: Type of benefit to extract ('latency' or 'cost') + benefit_type: 'latency', 'cost' (query CPU) or 'total_cost' (total CPU) Returns: DataFrame with experiment data @@ -393,8 +390,10 @@ def extract_experiments_from_patterns( for exp_name in sorted(exp_names): if benefit_type == "latency": exp_data = extract_experiment_data(exp_name, metric=metric) - elif benefit_type == "cost": - exp_data = extract_cost_benefit_data(exp_name, metric=metric) + elif benefit_type in ("cost", "total_cost"): + exp_data = extract_cost_benefit_data( + exp_name, metric=metric, total=benefit_type == "total_cost" + ) else: raise ValueError(f"Unknown benefit_type: {benefit_type}") @@ -463,7 +462,11 @@ def create_plot( y_breaks = sorted(list(set(y_breaks))) # Remove duplicates and sort # Dynamic Y-axis label based on benefit type - y_label = "Latency Benefit" if benefit_type == "latency" else "Cost Benefit (CPU)" + y_label = { + "latency": "Latency Benefit", + "cost": "Query CPU Benefit", + "total_cost": "Total CPU Benefit", + }[benefit_type] p = ( ggplot( @@ -510,7 +513,8 @@ def print_summary_table( header = f"Latency Benefit Analysis Summary ({metric.upper()} metric)" else: cpu_stat = METRIC_TO_CPU_STAT.get(metric, metric) - header = f"Cost Benefit Analysis Summary (CPU {cpu_stat.upper()})" + cpu_kind = "Total" if benefit_type == "total_cost" else "Query" + header = f"{cpu_kind} CPU Benefit Analysis Summary ({cpu_stat.upper()})" print("\n" + "=" * 100) print(header) @@ -598,8 +602,11 @@ def main(): "--benefit-type", type=str, default="latency", - choices=["latency", "cost"], - help="Type of benefit to plot: latency or cost (CPU) (default: latency)", + choices=["latency", "cost", "total_cost"], + help=( + "Type of benefit to plot: latency, cost (query CPU) or total_cost " + "(total CPU across all processes) (default: latency)" + ), ) parser.add_argument( "--cardinalities", @@ -627,7 +634,7 @@ def main(): parser.error("Must specify at least one of --print or --plot") # Validate metric compatibility with benefit type - if args.benefit_type == "cost" and args.metric == "mean": + if args.benefit_type != "latency" and args.metric == "mean": parser.error( "'mean' metric is not available for cost benefit. Use median, p95, p99, sum, or max" ) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py index 516ec089..d59b1678 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -19,6 +19,7 @@ import os import sys import re +import json import glob import yaml import argparse @@ -277,17 +278,18 @@ def extract_experiment_data( def extract_cost_benefit_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True + exp_name: str, metric: str = "p95", verify_scale: bool = True, total: bool = False ) -> Optional[Dict[str, Any]]: """ Extract cost benefit data from a single experiment. - Runs compare_costs.py and parses Query CPU Benefit output. + Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. Args: exp_name: Experiment name metric: CPU metric to use (median, p95, p99, sum, max) verify_scale: If True, verify data scale matches expected 2^card_exp + total: Use total CPU (all processes) instead of query CPU Returns: dict with experiment data or None if extraction fails @@ -331,6 +333,7 @@ def extract_cost_benefit_data( exp_name, "--all_experiment_modes", "--print", + "--machine-readable", ], capture_output=True, text=True, @@ -338,24 +341,16 @@ def extract_cost_benefit_data( cwd=script_dir, ) - # Parse output for Query CPU Benefit section - output = result.stdout + result.stderr + costs = json.loads(result.stdout) cpu_stat = METRIC_TO_CPU_STAT.get(metric, "p95") - - # Look for pattern: " : x" in Query CPU Benefit section - # Handle both numeric values and "inf" - pattern = rf"Query CPU Benefit.*?^\s+{cpu_stat}:\s+([\d.]+|inf)x" - match = re.search(pattern, output, re.MULTILINE | re.DOTALL) - - if not match: - print( - f"Warning: Could not find Query CPU Benefit '{cpu_stat}' for {exp_name}" - ) + benefits = ( + costs["benefit"]["cpu_percent"] if total else costs["query_cpu_benefit"] + ) + benefit_ratio = benefits[cpu_stat] + if np.isnan(benefit_ratio): # 0/0: both modes had zero cost + print(f"Warning: {cpu_stat} CPU benefit is NaN for {exp_name}") return None - value_str = match.group(1) - benefit_ratio = float(value_str) if value_str != "inf" else float("inf") - return { "experiment_name": exp_name, "query_type": metadata["query_type"], @@ -389,7 +384,7 @@ def extract_experiments_from_patterns( patterns: List of glob patterns for experiment names metric: Metric to use (latency or CPU stat depending on benefit_type) cardinalities: Optional list of cardinality exponents to include - benefit_type: Type of benefit to extract ('latency' or 'cost') + benefit_type: 'latency', 'cost' (query CPU) or 'total_cost' (total CPU) query_types: Optional list of query types to include (e.g., ['qot']) Returns: @@ -413,8 +408,10 @@ def extract_experiments_from_patterns( for exp_name in sorted(exp_names): if benefit_type == "latency": exp_data = extract_experiment_data(exp_name, metric=metric) - elif benefit_type == "cost": - exp_data = extract_cost_benefit_data(exp_name, metric=metric) + elif benefit_type in ("cost", "total_cost"): + exp_data = extract_cost_benefit_data( + exp_name, metric=metric, total=benefit_type == "total_cost" + ) else: raise ValueError(f"Unknown benefit_type: {benefit_type}") @@ -489,11 +486,11 @@ def create_plot( y_breaks = sorted(list(set(y_breaks))) # Remove duplicates and sort # Dynamic Y-axis label based on benefit type - y_label = ( - "Query Latency Benefit" - if benefit_type == "latency" - else "Query Cost Benefit (CPU)" - ) + y_label = { + "latency": "Query Latency Benefit", + "cost": "Query CPU Benefit", + "total_cost": "Total CPU Benefit", + }[benefit_type] p = ( ggplot( @@ -543,7 +540,8 @@ def print_summary_table( header = f"Latency Benefit Analysis Summary ({metric.upper()} metric)" else: cpu_stat = METRIC_TO_CPU_STAT.get(metric, metric) - header = f"Cost Benefit Analysis Summary (CPU {cpu_stat.upper()})" + cpu_kind = "Total" if benefit_type == "total_cost" else "Query" + header = f"{cpu_kind} CPU Benefit Analysis Summary ({cpu_stat.upper()})" print("\n" + "=" * 100) print(header) @@ -634,8 +632,11 @@ def main(): "--benefit-type", type=str, default="latency", - choices=["latency", "cost"], - help="Type of benefit to plot: latency or cost (CPU) (default: latency)", + choices=["latency", "cost", "total_cost"], + help=( + "Type of benefit to plot: latency, cost (query CPU) or total_cost " + "(total CPU across all processes) (default: latency)" + ), ) parser.add_argument( "--cardinalities", @@ -669,7 +670,7 @@ def main(): parser.error("Must specify at least one of --print or --plot") # Validate metric compatibility with benefit type - if args.benefit_type == "cost" and args.metric == "mean": + if args.benefit_type != "latency" and args.metric == "mean": parser.error( "'mean' metric is not available for cost benefit. Use median, p95, p99, sum, or max" ) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py index b9ab5cb0..f6379f44 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py @@ -47,8 +47,8 @@ def print_benefits_summary( cost_label = "Query Cost P95 Benefit (ratio)" cost_json_key = "query_cost_p95_benefit_ratio" else: - cost_label = "Cost P95 Benefit (ratio)" - cost_json_key = "cost_p95_benefit_ratio" + cost_label = "Total CPU P95 Benefit (ratio)" + cost_json_key = "total_cpu_p95_benefit_ratio" print("\nBenefits Summary (Prometheus / SketchDB):") print("=" * 110) @@ -137,14 +137,14 @@ def plot_scale_vs_benefits( # Determine cost label based on type if use_query_cost_sum: - cost_ylabel = "Query Cost Benefit (ratio)" - cost_legend = "Query Cost Benefit" + cost_ylabel = "Query CPU Sum Benefit (ratio)" + cost_legend = "Query CPU Sum Benefit" elif use_query_cost_95: - cost_ylabel = "Query Cost P95 Benefit (ratio)" - cost_legend = "Query Cost P95 Benefit" + cost_ylabel = "Query CPU P95 Benefit (ratio)" + cost_legend = "Query CPU P95 Benefit" else: - cost_ylabel = "Cost P95 Benefit (ratio)" - cost_legend = "Cost P95 Benefit" + cost_ylabel = "Total CPU P95 Benefit (ratio)" + cost_legend = "Total CPU P95 Benefit" # Plot cost benefit on left y-axis color_cost = "#1f77b4" @@ -280,12 +280,12 @@ def main(): parser.add_argument( "--use-query-cost-sum", action="store_true", - help="Use query CPU cost sum instead of p95 CPU cost", + help="Use query CPU sum instead of total CPU p95", ) parser.add_argument( "--use-query-cost-95", action="store_true", - help="Use query CPU cost p95 instead of p95 CPU cost", + help="Use query CPU p95 instead of total CPU p95", ) args = parser.parse_args() @@ -341,7 +341,7 @@ def main(): else: cost_prometheus = get_cost_p95(exp_name, "baseline") cost_sketchdb = get_cost_p95(exp_name, "sketchdb") - cost_type = "Cost p95" + cost_type = "Total CPU p95" if cost_prometheus is not None and cost_sketchdb is not None: cost_benefit = cost_prometheus / cost_sketchdb diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index 5526ee04..07f76028 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -8,7 +8,6 @@ import argparse import os import sys -import re import json import subprocess import yaml @@ -86,278 +85,96 @@ def calculate_data_scale(experiment_name): return None -def get_latency_p95(experiment_name, experiment_mode="baseline"): - """ - Get p95 latency by running ./run_compare_latencies.sh - - Args: - experiment_name: Name of the experiment - experiment_mode: Experiment mode (baseline or sketchdb) +def _run_json(script, args): + """Run a single_experiment script with --machine-readable and parse its JSON. - Returns: - p95 latency value (exact), or None if failed + Returns None if the script fails (e.g. missing experiment); a missing JSON + key raises instead, so a schema change can't silently drop data. """ - script_dir = SINGLE_EXPERIMENT_DIR - script_path = os.path.join(script_dir, "run_compare_latencies.sh") - try: result = subprocess.run( - [script_path, experiment_name], + ["python3", os.path.join(SINGLE_EXPERIMENT_DIR, script)] + + args + + ["--machine-readable"], capture_output=True, text=True, check=True, - cwd=script_dir, ) - - # Parse output to extract p95 from the appropriate line based on experiment mode - # For baseline: look for "exact:" line - # For sketchdb: look for "estimate:" line - output = result.stdout + result.stderr - - # Map experiment mode to the line to search for - search_key = "exact" if experiment_mode == "baseline" else "estimate" - - # Find the line with the appropriate key - pattern = rf"{search_key}:\s*\{{([^}}]+)\}}" - match = re.search(pattern, output) - - if match: - dict_str = match.group(1) - # Extract p95 value - p95_match = re.search(r"'p95':\s*([\d.]+)", dict_str) - if p95_match: - return float(p95_match.group(1)) - - print( - f"Warning: Could not parse p95 latency from output for {experiment_name} in {experiment_mode} mode" - ) - return None - except subprocess.CalledProcessError as e: - print(f"Error running latency comparison for {experiment_name}: {e}") + print(f"Error running {script} {args}: {e.stderr.strip()[-300:]}") return None - except Exception as e: - print(f"Error getting latency for {experiment_name}: {e}") - return None - - -def get_cost_p95(experiment_name, experiment_mode="baseline"): - """ - Get p95 CPU cost by running compare_costs.py - - Args: - experiment_name: Name of the experiment - experiment_mode: Experiment mode (baseline or sketchdb) - - Returns: - p95 CPU percentage, or None if failed - """ - script_dir = SINGLE_EXPERIMENT_DIR - compare_costs_path = os.path.join(script_dir, "compare_costs.py") - - try: - result = subprocess.run( - [ - "python3", - compare_costs_path, - "--experiment_name", - experiment_name, - "--experiment_mode", - experiment_mode, - "--print", - ], - capture_output=True, - text=True, - check=True, - cwd=script_dir, - ) + return json.loads(result.stdout) + + +def _compare_costs(experiment_name, experiment_mode): + return _run_json( + "compare_costs.py", + [ + "--experiment_name", + experiment_name, + "--experiment_mode", + experiment_mode, + "--print", + ], + ) - # Parse output to extract p95 CPU from "{experiment_mode} {experiment_mode}.yml cpu_percent p95" - # or "{experiment_mode} {experiment_mode} cpu_percent p95" - output = result.stdout + result.stderr - # Look for lines matching the pattern - pattern = ( - rf"{experiment_mode}\s+{experiment_mode}.*cpu_percent\s+p95\s+([\d.]+)" - ) - for line in output.split("\n"): - if re.search(pattern, line): - match = re.search(pattern, line) - if match: - return float(match.group(1)) - - print( - f"Warning: Could not parse p95 CPU cost from output for {experiment_name}" - ) +def get_latency_p95(experiment_name, experiment_mode="baseline"): + """p95 latency pooled over all queries (baseline = exact, sketchdb = estimate).""" + data = _run_json( + "compare_latencies.py", + [ + "--experiment_name", + experiment_name, + "--exact_experiment_mode", + "baseline", + "--estimate_experiment_mode", + "sketchdb", + ], + ) + if data is None: return None + side = "exact" if experiment_mode == "baseline" else "estimate" + return data["results"]["-1"][side]["p95"] - except subprocess.CalledProcessError as e: - print(f"Error running cost comparison for {experiment_name}: {e}") - return None - except Exception as e: - print(f"Error getting cost for {experiment_name}: {e}") + +def get_cost_p95(experiment_name, experiment_mode="baseline"): + """p95 of total CPU % (sum over all monitored processes: ingest + query).""" + data = _compare_costs(experiment_name, experiment_mode) + if data is None: return None + return data["experiment_modes"][experiment_mode]["processes"]["all_all"][ + "cpu_percent" + ]["p95"] def get_query_cost_95(experiment_name, experiment_mode="baseline"): - """ - Get query CPU cost p95 by running compare_costs.py - - Args: - experiment_name: Name of the experiment - experiment_mode: Experiment mode (baseline or sketchdb) - - Returns: - Query CPU cost p95 percentage, or None if failed - """ - script_dir = SINGLE_EXPERIMENT_DIR - compare_costs_path = os.path.join(script_dir, "compare_costs.py") - - try: - result = subprocess.run( - [ - "python3", - compare_costs_path, - "--experiment_name", - experiment_name, - "--experiment_mode", - experiment_mode, - "--print", - ], - capture_output=True, - text=True, - check=True, - cwd=script_dir, - ) - - # Parse output to extract query CPU sum from "Query CPU Statistics" section - # Looking for pattern like: - # {experiment_mode}: - # p95: 1122837.55% - output = result.stdout + result.stderr - - # Look for the Query CPU Statistics section - in_query_section = False - in_experiment_mode_subsection = False - for line in output.split("\n"): - if "Query CPU Statistics" in line: - in_query_section = True - continue - - if in_query_section: - # Check if we're in the experiment_mode subsection - if line.strip().startswith(f"{experiment_mode}:"): - in_experiment_mode_subsection = True - continue - - # If we're in experiment_mode subsection, look for sum - if in_experiment_mode_subsection: - match = re.search(r"p95:\s+([\d.]+)%", line) - if match: - return float(match.group(1)) - # If we hit another section, stop - if line.strip() and not line.strip().startswith( - ("sum:", "max:", "median:", "p95:", "p99:") - ): - break - - print( - f"Warning: Could not parse query CPU cost sum from output for {experiment_name}" - ) - return None - - except subprocess.CalledProcessError as e: - print(f"Error running cost comparison for {experiment_name}: {e}") - return None - except Exception as e: - print(f"Error getting query cost sum for {experiment_name}: {e}") + """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" + data = _compare_costs(experiment_name, experiment_mode) + if data is None: return None + return data["query_cpu"][experiment_mode]["p95"] def get_query_cost_sum(experiment_name, experiment_mode="baseline"): - """ - Get query CPU cost sum by running compare_costs.py - - Args: - experiment_name: Name of the experiment - experiment_mode: Experiment mode (baseline or sketchdb) - - Returns: - Query CPU cost sum percentage, or None if failed - """ - script_dir = SINGLE_EXPERIMENT_DIR - compare_costs_path = os.path.join(script_dir, "compare_costs.py") - - try: - result = subprocess.run( - [ - "python3", - compare_costs_path, - "--experiment_name", - experiment_name, - "--experiment_mode", - experiment_mode, - "--print", - ], - capture_output=True, - text=True, - check=True, - cwd=script_dir, - ) - - # Parse output to extract query CPU sum from "Query CPU Statistics" section - # Looking for pattern like: - # {experiment_mode}: - # sum: 1122837.55% - output = result.stdout + result.stderr - - # Look for the Query CPU Statistics section - in_query_section = False - in_experiment_mode_subsection = False - for line in output.split("\n"): - if "Query CPU Statistics" in line: - in_query_section = True - continue - - if in_query_section: - # Check if we're in the experiment_mode subsection - if line.strip().startswith(f"{experiment_mode}:"): - in_experiment_mode_subsection = True - continue - - # If we're in experiment_mode subsection, look for sum - if in_experiment_mode_subsection: - match = re.search(r"sum:\s+([\d.]+)%", line) - if match: - return float(match.group(1)) - # If we hit another section, stop - if line.strip() and not line.strip().startswith( - ("sum:", "max:", "median:", "p95:", "p99:") - ): - break - - print( - f"Warning: Could not parse query CPU cost sum from output for {experiment_name}" - ) + """Sum of query CPU % over the run; depends on run length.""" + data = _compare_costs(experiment_name, experiment_mode) + if data is None: return None + return data["query_cpu"][experiment_mode]["sum"] - except subprocess.CalledProcessError as e: - print(f"Error running cost comparison for {experiment_name}: {e}") - return None - except Exception as e: - print(f"Error getting query cost sum for {experiment_name}: {e}") - return None +def cost_label_for(use_query_cost_sum=False, use_query_cost_95=False): + if use_query_cost_sum: + return "Query CPU sum (%)" + if use_query_cost_95: + return "Query CPU p95 (%)" + return "Total CPU p95 (%)" -def print_data_summary( - experiments, data_scales, latencies, costs, use_query_cost_sum=False -): - """Print summary of the data.""" - cost_label = "Query Cost Sum (CPU %)" if use_query_cost_sum else "Cost P95 (CPU %)" - cost_json_key = ( - "query_cost_sum_cpu_percent" if use_query_cost_sum else "cost_p95_cpu_percent" - ) +def print_data_summary(experiments, data_scales, latencies, costs, cost_label): + """Print summary of the data.""" + cost_json_key = cost_label print("\nData Summary:") print("=" * 100) print( @@ -394,9 +211,9 @@ def plot_scale_vs_metrics( data_scales, latencies, costs, + cost_label, save_file=None, show=False, - use_query_cost_sum=False, ): """ Plot data scale vs cost and latency. @@ -405,10 +222,10 @@ def plot_scale_vs_metrics( experiments: List of experiment names data_scales: List of data scale values (metrics/sec) latencies: List of p95 latency values (seconds) - costs: List of CPU cost values (% - either p95 or query sum) + costs: List of CPU cost values (%) + cost_label: Axis label naming the cost definition (see cost_label_for) save_file: Filename to save the plot (if None, doesn't save) show: Whether to display the plot - use_query_cost_sum: Whether cost values represent query cost sum instead of p95 Returns: matplotlib figure object @@ -438,11 +255,7 @@ def plot_scale_vs_metrics( # Create the plot with two y-axes fig, ax1 = plt.subplots(figsize=(12, 6)) - # Determine cost label based on type - cost_ylabel = ( - "Query Cost (CPU %, sum)" if use_query_cost_sum else "p95 CPU usage (%)" - ) - cost_legend = "Query Cost (CPU %)" if use_query_cost_sum else "p95 CPU Usage (%)" + cost_ylabel = cost_legend = cost_label # Plot cost on left y-axis color_cost = "#1f77b4" @@ -546,12 +359,12 @@ def main(): parser.add_argument( "--use-query-cost-sum", action="store_true", - help="Use query CPU cost sum instead of p95 CPU cost", + help="Use query CPU sum instead of total CPU p95", ) parser.add_argument( "--use-query-cost-95", action="store_true", - help="Use query CPU cost p95 instead of p95 CPU cost", + help="Use query CPU p95 instead of total CPU p95", ) parser.add_argument( "--experiment_mode", @@ -562,6 +375,7 @@ def main(): ) args = parser.parse_args() + cost_label = cost_label_for(args.use_query_cost_sum, args.use_query_cost_95) # Validate arguments if args.plot and not (args.save or args.show): @@ -592,30 +406,19 @@ def main(): if latency is not None: print(f" Latency p95: {latency:.4f} seconds") - # Get cost (either p95 or query cost sum based on flag) if args.use_query_cost_sum: cost = get_query_cost_sum(exp_name, args.experiment_mode) - if cost is not None: - print(f" Query cost sum: {cost:.2f} CPU %") elif args.use_query_cost_95: cost = get_query_cost_95(exp_name, args.experiment_mode) - if cost is not None: - print(f" Query cost p95: {cost:.2f} CPU %") else: cost = get_cost_p95(exp_name, args.experiment_mode) - if cost is not None: - print(f" Cost p95: {cost:.2f} CPU %") + if cost is not None: + print(f" {cost_label}: {cost:.2f}") costs.append(cost) # Print summary if requested if args.print: - print_data_summary( - EXPERIMENT_NAMES, - data_scales, - latencies, - costs, - use_query_cost_sum=args.use_query_cost_sum, - ) + print_data_summary(EXPERIMENT_NAMES, data_scales, latencies, costs, cost_label) # Generate plot if requested if args.plot: @@ -624,9 +427,9 @@ def main(): data_scales=data_scales, latencies=latencies, costs=costs, + cost_label=cost_label, save_file=args.save, show=args.show, - use_query_cost_sum=args.use_query_cost_sum, ) return 0 From b3d37fc07b12eaf1552bd25879fe678fc338705f Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 18:09:17 -0400 Subject: [PATCH 3/6] fix(tools): label post_experiment latency in seconds, not ms Query latency is recorded as time.time() differences (seconds) and is never converted, but plot_latency_cost_tradeoff.py labeled it [ms]. Co-Authored-By: Claude Opus 5.5 --- .../multi_experiment/plot_latency_cost_tradeoff.py | 6 +++--- .../multi_experiment/plot_scale_vs_metrics.py | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py index 5c051116..e9dfca62 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py @@ -259,7 +259,7 @@ def plot_latency_cost_tradeoff( ), ) - ax.set_xlabel(f"Latency ({latency_metric}) [ms]") + ax.set_xlabel(f"Latency ({latency_metric}) [s]") ax.set_ylabel(f"{args.cpu_type.capitalize()} CPU Cost ({cost_metric}) [%]") ax.set_title("Latency-Cost Tradeoff: Prometheus vs ASAPOlly") ax.legend() @@ -322,10 +322,10 @@ def main(args): exact_cost, est_cost = exact_total, est_total data_points[exp_name] = (exact_lat, exact_cost, est_lat, est_cost) print( - f" Prometheus: latency={exact_lat:.2f}ms, {args.cpu_type}_cpu={exact_cost:.2f}%" + f" Prometheus: latency={exact_lat:.2f}s, {args.cpu_type}_cpu={exact_cost:.2f}%" ) print( - f" TurboProm: latency={est_lat:.2f}ms, {args.cpu_type}_cpu={est_cost:.2f}%" + f" TurboProm: latency={est_lat:.2f}s, {args.cpu_type}_cpu={est_cost:.2f}%" ) except Exception as e: print(f" Failed to process {exp_name}: {e}") diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index 07f76028..a04c53e7 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -2,7 +2,7 @@ """ Script to plot data scale vs cost and latency across multiple experiments. X-axis: Data scale (metrics/sec) in log scale -Y-axes: Left = Cost (CPU %), Right = Latency (ms) +Y-axes: Left = Cost (CPU %), Right = Latency (s) """ import argparse From 34be6a7fc33ce160e2476ae6fdae30cefdaebe9a Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 18:49:44 -0400 Subject: [PATCH 4/6] fix(tools): harden post_experiment plots against missing and infinite values - plot_latency_cost_tradeoff: only require the CPU stats for the chosen --cpu_type; drop unsupported --cost_metric mean. - plot_scale_vs_metrics: return None with a warning when compare_costs omits query CPU instead of raising KeyError. - plot_cardinality_vs_benefit(_v2): drop non-finite benefit ratios before plotting; inf made the y-tick loop unbounded. Co-Authored-By: Claude Opus 5.5 --- .../plot_cardinality_vs_benefit.py | 10 +++ .../plot_cardinality_vs_benefit_v2.py | 10 +++ .../plot_latency_cost_tradeoff.py | 65 ++++++------------- .../multi_experiment/plot_scale_vs_metrics.py | 20 ++++-- 4 files changed, 53 insertions(+), 52 deletions(-) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py index b8020e46..2ebedaaa 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py @@ -430,6 +430,16 @@ def create_plot( Returns: plotnine ggplot object """ + # inf (zero sketchdb cost) can't be placed on the axis and would make the + # tick loop below unbounded. + finite = np.isfinite(df["benefit_ratio"]) + if not finite.all(): + print( + "Warning: dropping non-finite benefit ratios from plot: " + f"{df.loc[~finite, 'experiment_name'].tolist()}" + ) + df = df[finite] + # Create labels for legend card_exps = sorted(df["card_exp"].unique()) color_labels = [f"2^{exp}" for exp in card_exps] diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py index d59b1678..d489fb85 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -452,6 +452,16 @@ def create_plot( Returns: plotnine ggplot object """ + # inf (zero sketchdb cost) can't be placed on the axis and would make the + # tick loop below unbounded. + finite = np.isfinite(df["benefit_ratio"]) + if not finite.all(): + print( + "Warning: dropping non-finite benefit ratios from plot: " + f"{df.loc[~finite, 'experiment_name'].tolist()}" + ) + df = df[finite] + # Create labels for legend card_exps = sorted(df["card_exp"].unique()) color_labels = [f"2^{exp}" for exp in card_exps] diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py index e9dfca62..d9b0727a 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py @@ -71,39 +71,29 @@ def extract_metrics( cost_metric: str, exact_mode: str, estimate_mode: str, -) -> Tuple[float, float, float, float, float, float]: + cpu_type: str, +) -> Tuple[float, float, float, float]: """Extract latency and cost metrics for both modes. + cpu_type "query" uses query-attributed CPU; "total" uses CPU summed over + all monitored processes ("all" pseudo-process in compare_costs.py). + Returns: - (exact_latency, exact_cost, estimate_latency, estimate_cost, - exact_total_cpu, estimate_total_cpu) + (exact_latency, exact_cost, estimate_latency, estimate_cost) """ - # Get cost data cost_data = run_compare_costs(experiment_name) - if "query_cpu" not in cost_data: - raise ValueError(f"No query_cpu data found for {experiment_name}") - - if exact_mode not in cost_data["query_cpu"]: - raise ValueError( - f"Mode {exact_mode} not found in query_cpu data for {experiment_name}" - ) - if estimate_mode not in cost_data["query_cpu"]: - raise ValueError( - f"Mode {estimate_mode} not found in query_cpu data for {experiment_name}" - ) + def cpu(mode): + if cpu_type == "total": + return cost_data["experiment_modes"][mode]["processes"]["all_all"][ + "cpu_percent" + ][cost_metric] + if mode not in cost_data.get("query_cpu", {}): + raise ValueError(f"No query_cpu data for mode {mode} in {experiment_name}") + return cost_data["query_cpu"][mode][cost_metric] - exact_cost = cost_data["query_cpu"][exact_mode][cost_metric] - estimate_cost = cost_data["query_cpu"][estimate_mode][cost_metric] - - # Total CPU across all monitored processes ("all" pseudo-process in compare_costs.py) - def total_cpu(mode): - return cost_data["experiment_modes"][mode]["processes"]["all_all"][ - "cpu_percent" - ][cost_metric] - - exact_total_cpu = total_cpu(exact_mode) - estimate_total_cpu = total_cpu(estimate_mode) + exact_cost = cpu(exact_mode) + estimate_cost = cpu(estimate_mode) # Get latency data latency_data = run_compare_latencies(experiment_name, exact_mode, estimate_mode) @@ -118,14 +108,7 @@ def total_cpu(mode): exact_latency = latency_data["results"]["-1"]["exact"][latency_metric] estimate_latency = latency_data["results"]["-1"]["estimate"][latency_metric] - return ( - exact_latency, - exact_cost, - estimate_latency, - estimate_cost, - exact_total_cpu, - estimate_total_cpu, - ) + return exact_latency, exact_cost, estimate_latency, estimate_cost def plot_latency_cost_tradeoff( @@ -304,22 +287,14 @@ def main(args): for exp_name in experiment_names: try: print(f"\nProcessing experiment: {exp_name}") - ( - exact_lat, - exact_cost, - est_lat, - est_cost, - exact_total, - est_total, - ) = extract_metrics( + exact_lat, exact_cost, est_lat, est_cost = extract_metrics( exp_name, args.latency_metric, args.cost_metric, args.exact_mode, args.estimate_mode, + args.cpu_type, ) - if args.cpu_type == "total": - exact_cost, est_cost = exact_total, est_total data_points[exp_name] = (exact_lat, exact_cost, est_lat, est_cost) print( f" Prometheus: latency={exact_lat:.2f}s, {args.cpu_type}_cpu={exact_cost:.2f}%" @@ -364,7 +339,7 @@ def main(args): "--cost_metric", type=str, default="p99", - choices=["median", "mean", "p95", "p99", "sum"], + choices=["median", "max", "p95", "p99", "sum"], help="Cost metric to use (default: p99)", ) parser.add_argument( diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index a04c53e7..a6976cc3 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -148,20 +148,26 @@ def get_cost_p95(experiment_name, experiment_mode="baseline"): ]["p95"] -def get_query_cost_95(experiment_name, experiment_mode="baseline"): - """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" +def _query_cpu(experiment_name, experiment_mode, stat): data = _compare_costs(experiment_name, experiment_mode) if data is None: return None - return data["query_cpu"][experiment_mode]["p95"] + # compare_costs omits query_cpu when it can't attribute CPU (e.g. no + # prometheus process in baseline). + if experiment_mode not in data.get("query_cpu", {}): + print(f"Warning: no query CPU for {experiment_name} ({experiment_mode})") + return None + return data["query_cpu"][experiment_mode][stat] + + +def get_query_cost_95(experiment_name, experiment_mode="baseline"): + """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" + return _query_cpu(experiment_name, experiment_mode, "p95") def get_query_cost_sum(experiment_name, experiment_mode="baseline"): """Sum of query CPU % over the run; depends on run length.""" - data = _compare_costs(experiment_name, experiment_mode) - if data is None: - return None - return data["query_cpu"][experiment_mode]["sum"] + return _query_cpu(experiment_name, experiment_mode, "sum") def cost_label_for(use_query_cost_sum=False, use_query_cost_95=False): From 86ea8eeb54b0214c5601166b574c8763de01c8f7 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 19:07:31 -0400 Subject: [PATCH 5/6] Revert "fix(tools): harden post_experiment plots against missing and infinite values" This reverts commit 34be6a7fc33ce160e2476ae6fdae30cefdaebe9a. --- .../plot_cardinality_vs_benefit.py | 10 --- .../plot_cardinality_vs_benefit_v2.py | 10 --- .../plot_latency_cost_tradeoff.py | 65 +++++++++++++------ .../multi_experiment/plot_scale_vs_metrics.py | 20 ++---- 4 files changed, 52 insertions(+), 53 deletions(-) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py index 2ebedaaa..b8020e46 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py @@ -430,16 +430,6 @@ def create_plot( Returns: plotnine ggplot object """ - # inf (zero sketchdb cost) can't be placed on the axis and would make the - # tick loop below unbounded. - finite = np.isfinite(df["benefit_ratio"]) - if not finite.all(): - print( - "Warning: dropping non-finite benefit ratios from plot: " - f"{df.loc[~finite, 'experiment_name'].tolist()}" - ) - df = df[finite] - # Create labels for legend card_exps = sorted(df["card_exp"].unique()) color_labels = [f"2^{exp}" for exp in card_exps] diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py index d489fb85..d59b1678 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -452,16 +452,6 @@ def create_plot( Returns: plotnine ggplot object """ - # inf (zero sketchdb cost) can't be placed on the axis and would make the - # tick loop below unbounded. - finite = np.isfinite(df["benefit_ratio"]) - if not finite.all(): - print( - "Warning: dropping non-finite benefit ratios from plot: " - f"{df.loc[~finite, 'experiment_name'].tolist()}" - ) - df = df[finite] - # Create labels for legend card_exps = sorted(df["card_exp"].unique()) color_labels = [f"2^{exp}" for exp in card_exps] diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py index d9b0727a..e9dfca62 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_cost_tradeoff.py @@ -71,29 +71,39 @@ def extract_metrics( cost_metric: str, exact_mode: str, estimate_mode: str, - cpu_type: str, -) -> Tuple[float, float, float, float]: +) -> Tuple[float, float, float, float, float, float]: """Extract latency and cost metrics for both modes. - cpu_type "query" uses query-attributed CPU; "total" uses CPU summed over - all monitored processes ("all" pseudo-process in compare_costs.py). - Returns: - (exact_latency, exact_cost, estimate_latency, estimate_cost) + (exact_latency, exact_cost, estimate_latency, estimate_cost, + exact_total_cpu, estimate_total_cpu) """ + # Get cost data cost_data = run_compare_costs(experiment_name) - def cpu(mode): - if cpu_type == "total": - return cost_data["experiment_modes"][mode]["processes"]["all_all"][ - "cpu_percent" - ][cost_metric] - if mode not in cost_data.get("query_cpu", {}): - raise ValueError(f"No query_cpu data for mode {mode} in {experiment_name}") - return cost_data["query_cpu"][mode][cost_metric] + if "query_cpu" not in cost_data: + raise ValueError(f"No query_cpu data found for {experiment_name}") + + if exact_mode not in cost_data["query_cpu"]: + raise ValueError( + f"Mode {exact_mode} not found in query_cpu data for {experiment_name}" + ) + if estimate_mode not in cost_data["query_cpu"]: + raise ValueError( + f"Mode {estimate_mode} not found in query_cpu data for {experiment_name}" + ) - exact_cost = cpu(exact_mode) - estimate_cost = cpu(estimate_mode) + exact_cost = cost_data["query_cpu"][exact_mode][cost_metric] + estimate_cost = cost_data["query_cpu"][estimate_mode][cost_metric] + + # Total CPU across all monitored processes ("all" pseudo-process in compare_costs.py) + def total_cpu(mode): + return cost_data["experiment_modes"][mode]["processes"]["all_all"][ + "cpu_percent" + ][cost_metric] + + exact_total_cpu = total_cpu(exact_mode) + estimate_total_cpu = total_cpu(estimate_mode) # Get latency data latency_data = run_compare_latencies(experiment_name, exact_mode, estimate_mode) @@ -108,7 +118,14 @@ def cpu(mode): exact_latency = latency_data["results"]["-1"]["exact"][latency_metric] estimate_latency = latency_data["results"]["-1"]["estimate"][latency_metric] - return exact_latency, exact_cost, estimate_latency, estimate_cost + return ( + exact_latency, + exact_cost, + estimate_latency, + estimate_cost, + exact_total_cpu, + estimate_total_cpu, + ) def plot_latency_cost_tradeoff( @@ -287,14 +304,22 @@ def main(args): for exp_name in experiment_names: try: print(f"\nProcessing experiment: {exp_name}") - exact_lat, exact_cost, est_lat, est_cost = extract_metrics( + ( + exact_lat, + exact_cost, + est_lat, + est_cost, + exact_total, + est_total, + ) = extract_metrics( exp_name, args.latency_metric, args.cost_metric, args.exact_mode, args.estimate_mode, - args.cpu_type, ) + if args.cpu_type == "total": + exact_cost, est_cost = exact_total, est_total data_points[exp_name] = (exact_lat, exact_cost, est_lat, est_cost) print( f" Prometheus: latency={exact_lat:.2f}s, {args.cpu_type}_cpu={exact_cost:.2f}%" @@ -339,7 +364,7 @@ def main(args): "--cost_metric", type=str, default="p99", - choices=["median", "max", "p95", "p99", "sum"], + choices=["median", "mean", "p95", "p99", "sum"], help="Cost metric to use (default: p99)", ) parser.add_argument( diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index a6976cc3..a04c53e7 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -148,26 +148,20 @@ def get_cost_p95(experiment_name, experiment_mode="baseline"): ]["p95"] -def _query_cpu(experiment_name, experiment_mode, stat): +def get_query_cost_95(experiment_name, experiment_mode="baseline"): + """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" data = _compare_costs(experiment_name, experiment_mode) if data is None: return None - # compare_costs omits query_cpu when it can't attribute CPU (e.g. no - # prometheus process in baseline). - if experiment_mode not in data.get("query_cpu", {}): - print(f"Warning: no query CPU for {experiment_name} ({experiment_mode})") - return None - return data["query_cpu"][experiment_mode][stat] - - -def get_query_cost_95(experiment_name, experiment_mode="baseline"): - """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" - return _query_cpu(experiment_name, experiment_mode, "p95") + return data["query_cpu"][experiment_mode]["p95"] def get_query_cost_sum(experiment_name, experiment_mode="baseline"): """Sum of query CPU % over the run; depends on run length.""" - return _query_cpu(experiment_name, experiment_mode, "sum") + data = _compare_costs(experiment_name, experiment_mode) + if data is None: + return None + return data["query_cpu"][experiment_mode]["sum"] def cost_label_for(use_query_cost_sum=False, use_query_cost_95=False): From b859b77a69210231ec14c726ce5bfc983d6cd0e5 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Thu, 24 Sep 2026 19:14:13 -0400 Subject: [PATCH 6/6] refactor(tools): drop unneeded defaults in post_experiment scripts Every caller already passes these arguments; removing the defaults makes a forgotten argument (e.g. experiment_mode, metric, benefit_type, total) an error instead of silently selecting baseline/p95/query CPU. Inline verify_scale (always True), analyze_throughput's label_filter (always the float-only filter) and calculate_stable_throughput's num_windows (always self.num_windows). Co-Authored-By: Claude Opus 5.5 --- .../plot_cardinality_vs_benefit.py | 46 ++++++++---------- .../plot_cardinality_vs_benefit_v2.py | 48 ++++++++----------- .../multi_experiment/plot_latency_metrics.py | 4 +- .../plot_scale_vs_benefits.py | 8 ++-- .../multi_experiment/plot_scale_vs_metrics.py | 10 ++-- .../single_experiment/analyze_throughput.py | 41 +++++++--------- .../single_experiment/compare_costs.py | 4 +- 7 files changed, 68 insertions(+), 93 deletions(-) diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py index b8020e46..e5bece54 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit.py @@ -152,16 +152,15 @@ def load_experiment_config(exp_dir: str) -> Dict[str, Any]: return yaml.safe_load(f) -def extract_experiment_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True -) -> Optional[Dict[str, Any]]: +def extract_experiment_data(exp_name: str, metric: str) -> Optional[Dict[str, Any]]: """ Extract data from a single experiment. + Warns if the data scale doesn't match the expected 2^card_exp. + Args: exp_name: Experiment name metric: Latency metric to use (median, p95, p99, mean) - verify_scale: If True, verify data scale matches expected 2^card_exp Returns: dict with experiment data or None if extraction fails @@ -186,7 +185,7 @@ def extract_experiment_data( actual_scale = calculate_data_scale_from_config(config) expected_scale = 2 ** metadata["card_exp"] - if verify_scale and actual_scale != expected_scale: + if actual_scale != expected_scale: print( f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" ) @@ -262,17 +261,17 @@ def extract_experiment_data( def extract_cost_benefit_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True, total: bool = False + exp_name: str, metric: str, total: bool ) -> Optional[Dict[str, Any]]: """ Extract cost benefit data from a single experiment. - Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. + Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. Warns if + the data scale doesn't match the expected 2^card_exp. Args: exp_name: Experiment name metric: CPU metric to use (median, p95, p99, sum, max) - verify_scale: If True, verify data scale matches expected 2^card_exp total: Use total CPU (all processes) instead of query CPU Returns: @@ -291,16 +290,13 @@ def extract_cost_benefit_data( return None try: - # Verify scale if requested - actual_scale = None - if verify_scale: - config = load_experiment_config(exp_dir) - actual_scale = calculate_data_scale_from_config(config) - expected_scale = 2 ** metadata["card_exp"] - if actual_scale != expected_scale: - print( - f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" - ) + config = load_experiment_config(exp_dir) + actual_scale = calculate_data_scale_from_config(config) + expected_scale = 2 ** metadata["card_exp"] + if actual_scale != expected_scale: + print( + f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" + ) # Run compare_costs.py script_dir = os.path.join( @@ -356,9 +352,9 @@ def extract_cost_benefit_data( def extract_experiments_from_patterns( patterns: List[str], - metric: str = "p95", - cardinalities: Optional[List[int]] = None, - benefit_type: str = "latency", + metric: str, + cardinalities: Optional[List[int]], + benefit_type: str, ) -> pd.DataFrame: """ Extract data from experiments matching glob patterns. @@ -416,9 +412,7 @@ def extract_experiments_from_patterns( return df -def create_plot( - df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" -) -> "ggplot": +def create_plot(df: pd.DataFrame, metric: str, benefit_type: str) -> "ggplot": """ Create benefit vs lookback plot with log2(T/15) x-axis. @@ -504,9 +498,7 @@ def create_plot( return p -def print_summary_table( - df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" -): +def print_summary_table(df: pd.DataFrame, metric: str, benefit_type: str): """Print summary table of experiment data.""" # Dynamic header based on benefit type if benefit_type == "latency": diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py index d59b1678..46e80958 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -163,16 +163,15 @@ def load_experiment_config(exp_dir: str) -> Dict[str, Any]: return yaml.safe_load(f) -def extract_experiment_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True -) -> Optional[Dict[str, Any]]: +def extract_experiment_data(exp_name: str, metric: str) -> Optional[Dict[str, Any]]: """ Extract data from a single experiment. + Warns if the data scale doesn't match the expected 2^card_exp. + Args: exp_name: Experiment name metric: Latency metric to use (median, p95, p99, mean) - verify_scale: If True, verify data scale matches expected 2^card_exp Returns: dict with experiment data or None if extraction fails @@ -197,7 +196,7 @@ def extract_experiment_data( actual_scale = calculate_data_scale_from_config(config) expected_scale = 2 ** metadata["card_exp"] - if verify_scale and actual_scale != expected_scale: + if actual_scale != expected_scale: print( f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" ) @@ -278,17 +277,17 @@ def extract_experiment_data( def extract_cost_benefit_data( - exp_name: str, metric: str = "p95", verify_scale: bool = True, total: bool = False + exp_name: str, metric: str, total: bool ) -> Optional[Dict[str, Any]]: """ Extract cost benefit data from a single experiment. - Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. + Runs compare_costs.py and reads the baseline/sketchdb CPU ratio. Warns if + the data scale doesn't match the expected 2^card_exp. Args: exp_name: Experiment name metric: CPU metric to use (median, p95, p99, sum, max) - verify_scale: If True, verify data scale matches expected 2^card_exp total: Use total CPU (all processes) instead of query CPU Returns: @@ -307,16 +306,13 @@ def extract_cost_benefit_data( return None try: - # Verify scale if requested - actual_scale = None - if verify_scale: - config = load_experiment_config(exp_dir) - actual_scale = calculate_data_scale_from_config(config) - expected_scale = 2 ** metadata["card_exp"] - if actual_scale != expected_scale: - print( - f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" - ) + config = load_experiment_config(exp_dir) + actual_scale = calculate_data_scale_from_config(config) + expected_scale = 2 ** metadata["card_exp"] + if actual_scale != expected_scale: + print( + f"Warning: {exp_name} has scale {actual_scale} but expected {expected_scale}" + ) # Run compare_costs.py script_dir = os.path.join( @@ -372,10 +368,10 @@ def extract_cost_benefit_data( def extract_experiments_from_patterns( patterns: List[str], - metric: str = "p95", - cardinalities: Optional[List[int]] = None, - benefit_type: str = "latency", - query_types: Optional[List[str]] = None, + metric: str, + cardinalities: Optional[List[int]], + benefit_type: str, + query_types: Optional[List[str]], ) -> pd.DataFrame: """ Extract data from experiments matching glob patterns. @@ -438,9 +434,7 @@ def extract_experiments_from_patterns( return df -def create_plot( - df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" -) -> "ggplot": +def create_plot(df: pd.DataFrame, metric: str, benefit_type: str) -> "ggplot": """ Create benefit vs lookback plot with log2(T/15) x-axis. @@ -531,9 +525,7 @@ def create_plot( return p -def print_summary_table( - df: pd.DataFrame, metric: str = "p95", benefit_type: str = "latency" -): +def print_summary_table(df: pd.DataFrame, metric: str, benefit_type: str): """Print summary table of experiment data.""" # Dynamic header based on benefit type if benefit_type == "latency": diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py index 7b1b7c49..82fb94a3 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_latency_metrics.py @@ -151,8 +151,8 @@ class DataProcessor: def process_for_plotting( self, experiment_data: List[Dict[str, Any]], - individual_queries: bool = False, - show_benefit: bool = False, + individual_queries: bool, + show_benefit: bool, ) -> pd.DataFrame: """Process experiment data into format suitable for plotting.""" plot_data = [] diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py index f6379f44..3cb47c37 100644 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_benefits.py @@ -36,8 +36,8 @@ def print_benefits_summary( data_scales, latency_benefits, cost_benefits, - use_query_cost_sum=False, - use_query_cost_95=False, + use_query_cost_sum, + use_query_cost_95, ): """Print summary of the benefits data.""" if use_query_cost_sum: @@ -90,10 +90,10 @@ def plot_scale_vs_benefits( data_scales, latency_benefits, cost_benefits, + use_query_cost_sum, + use_query_cost_95, save_file=None, show=False, - use_query_cost_sum=False, - use_query_cost_95=False, ): """ Plot data scale vs benefits (prometheus/sketchdb ratios). diff --git a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py index a04c53e7..8aa744a2 100755 --- a/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_scale_vs_metrics.py @@ -119,7 +119,7 @@ def _compare_costs(experiment_name, experiment_mode): ) -def get_latency_p95(experiment_name, experiment_mode="baseline"): +def get_latency_p95(experiment_name, experiment_mode): """p95 latency pooled over all queries (baseline = exact, sketchdb = estimate).""" data = _run_json( "compare_latencies.py", @@ -138,7 +138,7 @@ def get_latency_p95(experiment_name, experiment_mode="baseline"): return data["results"]["-1"][side]["p95"] -def get_cost_p95(experiment_name, experiment_mode="baseline"): +def get_cost_p95(experiment_name, experiment_mode): """p95 of total CPU % (sum over all monitored processes: ingest + query).""" data = _compare_costs(experiment_name, experiment_mode) if data is None: @@ -148,7 +148,7 @@ def get_cost_p95(experiment_name, experiment_mode="baseline"): ]["p95"] -def get_query_cost_95(experiment_name, experiment_mode="baseline"): +def get_query_cost_95(experiment_name, experiment_mode): """p95 of query CPU % (see compare_costs.calculate_query_cpu).""" data = _compare_costs(experiment_name, experiment_mode) if data is None: @@ -156,7 +156,7 @@ def get_query_cost_95(experiment_name, experiment_mode="baseline"): return data["query_cpu"][experiment_mode]["p95"] -def get_query_cost_sum(experiment_name, experiment_mode="baseline"): +def get_query_cost_sum(experiment_name, experiment_mode): """Sum of query CPU % over the run; depends on run length.""" data = _compare_costs(experiment_name, experiment_mode) if data is None: @@ -164,7 +164,7 @@ def get_query_cost_sum(experiment_name, experiment_mode="baseline"): return data["query_cpu"][experiment_mode]["sum"] -def cost_label_for(use_query_cost_sum=False, use_query_cost_95=False): +def cost_label_for(use_query_cost_sum, use_query_cost_95): if use_query_cost_sum: return "Query CPU sum (%)" if use_query_cost_95: diff --git a/asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py b/asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py index 527f9e61..54d3f792 100755 --- a/asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py +++ b/asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py @@ -21,7 +21,7 @@ class ThroughputAnalyzer: """Analyzes throughput from prometheus metrics.""" - def __init__(self, window_duration: int = 30, num_windows: int = 10): + def __init__(self, window_duration: int, num_windows: int): """ Initialize the throughput analyzer. @@ -45,27 +45,26 @@ def load_prometheus_metrics(self, file_path: Path) -> Dict: raise def extract_timeseries( - self, - data: Dict, - metric_name: str, - label_filter: Optional[Dict[str, str]] = None, + self, data: Dict, metric_name: str ) -> List[Tuple[float, float]]: """ Extract timeseries data from prometheus metrics. + For prometheus_tsdb_head_samples_appended_total, only float samples + (label type="float") are counted. + Args: data: Loaded prometheus metrics data metric_name: Name of the metric to extract - label_filter: Dict of label key-value pairs to filter on (e.g., {"type": "float"}) Returns: List of (timestamp_seconds, value) tuples, sorted by timestamp """ - if ( - label_filter is None - and metric_name == "prometheus_tsdb_head_samples_appended_total" - ): - label_filter = {"type": "float"} # Default to float type only + label_filter: Optional[Dict[str, str]] = ( + {"type": "float"} + if metric_name == "prometheus_tsdb_head_samples_appended_total" + else None + ) timeseries = [] collection_start = datetime.fromisoformat(data["collection_start"]) @@ -103,15 +102,15 @@ def extract_timeseries( def calculate_rates( self, timeseries: List[Tuple[float, float]], - window_duration: Optional[int] = None, + window_duration: Optional[int], ) -> List[Tuple[float, float]]: """ Calculate rate (samples/sec) between measurements. Args: timeseries: List of (timestamp, cumulative_value) tuples - window_duration: If provided, only calculate rates for pairs separated by - approximately this duration (in seconds) + window_duration: If None, rates between consecutive points; otherwise + only pairs separated by approximately this duration (in seconds) Returns: List of (timestamp, rate) tuples where timestamp is the end of the interval @@ -166,21 +165,17 @@ def calculate_rates( return rates - def calculate_stable_throughput( - self, rates: List[Tuple[float, float]], num_windows: Optional[int] = None - ) -> float: + def calculate_stable_throughput(self, rates: List[Tuple[float, float]]) -> float: """ - Calculate stable throughput by averaging the last N rate measurements. + Calculate stable throughput by averaging the last self.num_windows rates. Args: rates: List of (timestamp, rate) tuples - num_windows: Number of last measurements to average (defaults to self.num_windows) Returns: Average rate over the last num_windows measurements """ - if num_windows is None: - num_windows = self.num_windows + num_windows = self.num_windows if len(rates) < num_windows: logger.warning( @@ -234,9 +229,7 @@ def analyze_prometheus(self, file_path: Path) -> Dict: ) # Calculate stable throughput - stable_throughput = self.calculate_stable_throughput( - windowed_rates, num_windows=self.num_windows - ) + stable_throughput = self.calculate_stable_throughput(windowed_rates) results[metric] = { "file": str(file_path), diff --git a/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py b/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py index b5f1c68f..abf126aa 100644 --- a/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py +++ b/asap-tools/experiments/post_experiment/single_experiment/compare_costs.py @@ -58,9 +58,7 @@ def measure_prometheus_ingest_cost(ingest_only_experiment_name: str) -> float: ) -def calculate_query_cpu( - monitor_info, experiment_mode, ingest_baseline_cpu_percent=None -): +def calculate_query_cpu(monitor_info, experiment_mode, ingest_baseline_cpu_percent): """ Calculate Query CPU timeseries for the given experiment mode.