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 87% 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..e5bece54 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 @@ -19,6 +19,7 @@ import os import sys import re +import json import glob import yaml import argparse @@ -45,13 +46,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 = { @@ -147,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 @@ -181,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}" ) @@ -257,17 +261,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, total: bool ) -> 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. 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: dict with experiment data or None if extraction fails @@ -285,19 +290,19 @@ 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.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( @@ -308,6 +313,7 @@ def extract_cost_benefit_data( exp_name, "--all_experiment_modes", "--print", + "--machine-readable", ], capture_output=True, text=True, @@ -315,22 +321,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"], @@ -352,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. @@ -363,7 +363,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 @@ -386,8 +386,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}") @@ -410,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. @@ -456,7 +456,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( @@ -494,16 +498,15 @@ 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": 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) @@ -591,8 +594,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", @@ -620,7 +626,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 new file mode 100644 index 00000000..46e80958 --- /dev/null +++ b/asap-tools/experiments/post_experiment/multi_experiment/plot_cardinality_vs_benefit_v2.py @@ -0,0 +1,702 @@ +#!/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 json +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) -> 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) + + 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 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, 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. 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) + total: Use total CPU (all processes) instead of query CPU + + 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: + 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", + "--machine-readable", + ], + capture_output=True, + text=True, + check=True, + cwd=script_dir, + ) + + costs = json.loads(result.stdout) + cpu_stat = METRIC_TO_CPU_STAT.get(metric, "p95") + 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 + + 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, + cardinalities: Optional[List[int]], + benefit_type: str, + query_types: Optional[List[str]], +) -> 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: 'latency', 'cost' (query CPU) or 'total_cost' (total CPU) + 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 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}") + + 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, benefit_type: str) -> "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 = { + "latency": "Query Latency Benefit", + "cost": "Query CPU Benefit", + "total_cost": "Total CPU Benefit", + }[benefit_type] + + 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, benefit_type: str): + """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) + 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) + 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", "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", + 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 != "latency" 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..e9dfca62 --- /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}) [s]") + 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}s, {args.cpu_type}_cpu={exact_cost:.2f}%" + ) + print( + 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}") + 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..82fb94a3 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: @@ -147,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 new file mode 100644 index 00000000..3cb47c37 --- /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, + use_query_cost_95, +): + """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 = "Total CPU P95 Benefit (ratio)" + cost_json_key = "total_cpu_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, + use_query_cost_sum, + use_query_cost_95, + save_file=None, + show=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 CPU Sum Benefit (ratio)" + cost_legend = "Query CPU Sum Benefit" + elif use_query_cost_95: + cost_ylabel = "Query CPU P95 Benefit (ratio)" + cost_legend = "Query CPU P95 Benefit" + else: + cost_ylabel = "Total CPU P95 Benefit (ratio)" + cost_legend = "Total CPU 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 sum instead of total CPU p95", + ) + parser.add_argument( + "--use-query-cost-95", + action="store_true", + help="Use query CPU p95 instead of total CPU p95", + ) + + 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 = "Total CPU 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 50% 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..8aa744a2 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 @@ -2,20 +2,22 @@ """ 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 import os import sys -import re import json import subprocess import yaml 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,267 +85,96 @@ def calculate_data_scale(experiment_name): return None -def get_latency_p95(experiment_name): - """ - Get p95 latency by running ./run_compare_latencies.sh - - Args: - experiment_name: Name of the experiment +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 = os.path.dirname(os.path.abspath(__file__)) - 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 exact - # Looking for: exact: {'median': X, 'p95': Y, ...} - 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) - # Extract p95 value - p95_match = re.search(r"'p95':\s*([\d.]+)", exact_dict_str) - if p95_match: - return float(p95_match.group(1)) - - print(f"Warning: Could not parse p95 latency from output for {experiment_name}") - return None - except subprocess.CalledProcessError as e: - print(f"Error running latency comparison for {experiment_name}: {e}") - return None - except Exception as e: - print(f"Error getting latency for {experiment_name}: {e}") + print(f"Error running {script} {args}: {e.stderr.strip()[-300:]}") return None + 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", + ], + ) -def get_cost_p95(experiment_name): - """ - Get p95 CPU cost by running compare_costs.py - - Args: - experiment_name: Name of the experiment - - Returns: - p95 CPU percentage, or None if failed - """ - script_dir = os.path.dirname(os.path.abspath(__file__)) - 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", - "baseline", - "--print", - ], - capture_output=True, - text=True, - check=True, - cwd=script_dir, - ) - - # Parse output to extract p95 CPU from "prometheus prometheus.yml cpu_percent p95" - # or "prometheus prometheus cpu_percent p95" - output = result.stdout + result.stderr - - # Look for lines matching the pattern - 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 match: - return float(match.group(1)) - - print( - f"Warning: Could not parse p95 CPU cost 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 cost for {experiment_name}: {e}") +def get_latency_p95(experiment_name, experiment_mode): + """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"] -def get_query_cost_95(experiment_name): - """ - Get query CPU cost p95 by running compare_costs.py - - Args: - experiment_name: Name of the experiment - - Returns: - Query CPU cost p95 percentage, or None if failed - """ - script_dir = os.path.dirname(os.path.abspath(__file__)) - 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", - "baseline", - "--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: - # prometheus: - # p95: 1122837.55% - output = result.stdout + result.stderr - - # Look for the Query CPU Statistics section - in_query_section = False - in_prometheus_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 - continue - - # If we're in prometheus subsection, look for sum - if in_prometheus_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}") +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: return None + return data["experiment_modes"][experiment_mode]["processes"]["all_all"][ + "cpu_percent" + ]["p95"] -def get_query_cost_sum(experiment_name): - """ - Get query CPU cost sum by running compare_costs.py - - Args: - experiment_name: Name of the experiment - - Returns: - Query CPU cost sum percentage, or None if failed - """ - script_dir = os.path.dirname(os.path.abspath(__file__)) - compare_costs_path = os.path.join(script_dir, "compare_costs.py") +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: + return None + return data["query_cpu"][experiment_mode]["p95"] - try: - result = subprocess.run( - [ - "python3", - compare_costs_path, - "--experiment_name", - experiment_name, - "--experiment_mode", - "baseline", - "--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: - # prometheus: - # sum: 1122837.55% - output = result.stdout + result.stderr - - # Look for the Query CPU Statistics section - in_query_section = False - in_prometheus_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 - continue - - # If we're in prometheus subsection, look for sum - if in_prometheus_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}" - ) +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: 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, use_query_cost_95): + 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( @@ -380,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. @@ -391,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 @@ -424,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" @@ -532,15 +359,23 @@ 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", + type=str, + choices=["baseline", "sketchdb"], + default="baseline", + help="Experiment mode (baseline or sketchdb)", ) 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): @@ -566,35 +401,24 @@ 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) - if cost is not None: - print(f" Query cost sum: {cost:.2f} CPU %") + cost = get_query_cost_sum(exp_name, args.experiment_mode) elif args.use_query_cost_95: - cost = get_query_cost_95(exp_name) - if cost is not None: - print(f" Query cost p95: {cost:.2f} CPU %") + cost = get_query_cost_95(exp_name, args.experiment_mode) else: - cost = get_cost_p95(exp_name) - if cost is not None: - print(f" Cost p95: {cost:.2f} CPU %") + cost = get_cost_p95(exp_name, args.experiment_mode) + 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: @@ -603,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 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 92% rename from asap-tools/experiments/post_experiment/analyze_throughput.py rename to asap-tools/experiments/post_experiment/single_experiment/analyze_throughput.py index 527f9e61..54d3f792 100755 --- a/asap-tools/experiments/post_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/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..abf126aa 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"] @@ -56,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. 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 \