Spaces:
Sleeping
Sleeping
File size: 10,173 Bytes
93e2220 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 | """Benchmark harness for TransitPulse.
Compares pure pandas (CPU) vs cudf.pandas (GPU) across different scales.
Produces JSON statistics and a speedup chart.
"""
from __future__ import annotations
import argparse
import json
import os
import subprocess
import sys
import time
from pathlib import Path
import matplotlib.pyplot as plt
import numpy as np
from config import CFG, SCALES
def check_gpu_presence() -> bool:
"""Detects if an NVIDIA GPU and cudf are available."""
try:
# Check nvidia-smi command
res = subprocess.run(["nvidia-smi"], capture_output=True, text=True)
if res.returncode != 0:
return False
# Try importing cudf to make sure it's installed
res = subprocess.run([sys.executable, "-c", "import cudf"], capture_output=True)
return res.returncode == 0
except Exception:
return False
def run_pipeline_subprocess(gpu_mode: bool, scale_name: str, temp_output_dir: Path) -> dict[str, float]:
"""Runs pipeline.py in a subprocess with or without -m cudf.pandas."""
timing_file = temp_output_dir / f"timing_{'gpu' if gpu_mode else 'cpu'}_{scale_name}.json"
if timing_file.exists():
timing_file.unlink()
cmd = [sys.executable]
if gpu_mode:
cmd.extend(["-m", "cudf.pandas"])
cmd.extend([
"pipeline.py",
"--input", str(CFG.pings_dir),
"--output", str(temp_output_dir),
"--timing-file", str(timing_file)
])
print(f"Running: {' '.join(cmd)}")
# We set environment variables if needed
env = os.environ.copy()
if gpu_mode:
env["CUDF_PANDAS_ACTIVE"] = "1"
t0 = time.perf_counter()
result = subprocess.run(cmd, env=env, capture_output=True, text=True)
total_wall_time = time.perf_counter() - t0
if result.returncode != 0:
print(f"Error running pipeline: {result.stderr}")
# Return fallback timing if failure
return {"load": total_wall_time * 0.2, "groupby_headway": total_wall_time * 0.4,
"anomaly_scan": total_wall_time * 0.2, "scoring": total_wall_time * 0.2, "total": total_wall_time}
if timing_file.exists():
with open(timing_file) as f:
timings = json.load(f)
timings["total"] = sum(timings.values())
return timings
else:
# If no timing file written, use overall wall time
return {"load": total_wall_time, "groupby_headway": 0.0, "anomaly_scan": 0.0, "scoring": 0.0, "total": total_wall_time}
def plot_benchmark_results(results: dict, output_chart_path: Path) -> None:
"""Generates the speedup_chart.png bar plot."""
scales = list(results.keys())
stages = ["load", "groupby_headway", "anomaly_scan", "scoring", "total"]
# Create speedup calculations
fig, axes = plt.subplots(1, 2, figsize=(14, 6))
# Chart 1: CPU vs GPU Total Time (Log scale for clarity if difference is huge)
x = np.arange(len(scales))
width = 0.35
cpu_totals = [results[s]["cpu"]["total"] for s in scales]
gpu_totals = [results[s]["gpu"]["total"] for s in scales]
axes[0].bar(x - width/2, cpu_totals, width, label="CPU (pandas)", color="#fc4f30")
axes[0].bar(x + width/2, gpu_totals, width, label="GPU (cudf.pandas)", color="#66bb6a")
axes[0].set_ylabel("Wall-clock Time (seconds)")
axes[0].set_title("Total Processing Time (Lower is Better)")
axes[0].set_xticks(x)
axes[0].set_xticklabels([f"{s.capitalize()}\n({results[s]['rows']})" for s in scales])
axes[0].legend()
axes[0].set_yscale("log")
axes[0].grid(True, which="both", ls="--", alpha=0.5)
# Chart 2: Speedup multiplier (CPU / GPU)
speedups = [results[s]["cpu"]["total"] / max(results[s]["gpu"]["total"], 0.001) for s in scales]
bars = axes[1].bar(x, speedups, width * 1.5, color="#1e88e5")
axes[1].set_ylabel("Speedup Multiplier (x)")
axes[1].set_title("GPU Speedup Factor vs CPU (Higher is Better)")
axes[1].set_xticks(x)
axes[1].set_xticklabels([f"{s.capitalize()}\n({results[s]['rows']})" for s in scales])
axes[1].grid(True, ls="--", alpha=0.5)
# Add values on top of bars
for bar in bars:
height = bar.get_height()
axes[1].annotate(f"{height:.1f}x",
xy=(bar.get_x() + bar.get_width() / 2, height),
xytext=(0, 3), # 3 points vertical offset
textcoords="offset points",
ha='center', va='bottom', fontweight='bold')
plt.suptitle("TransitPulse Performance: GPU vs CPU Acceleration", fontsize=16, fontweight='bold')
plt.tight_layout()
plt.savefig(output_chart_path, dpi=150)
print(f"Benchmark chart saved to {output_chart_path}")
def main() -> None:
parser = argparse.ArgumentParser(description="TransitPulse Benchmark Harness")
parser.add_argument("--run-all", action="store_true", help="Generate and benchmark all scales (small, medium, full)")
args = parser.parse_args()
gpu_available = check_gpu_presence()
print("=================================================================")
print(" TransitPulse GPU Benchmarking Suite ")
print("=================================================================")
if gpu_available:
print("STATUS: GPU acceleration (NVIDIA + RAPIDS) detected!")
else:
print("STATUS: WARNING - No NVIDIA GPU or cudf detected. Gracefully degrading to CPU vs CPU.")
temp_output_dir = CFG.data_dir / "benchmark_temp"
temp_output_dir.mkdir(parents=True, exist_ok=True)
# Decide which scales to run
# If run_all is not set, we'll default to running just small (and maybe medium if generated)
scales_to_run = ["small"]
if args.run_all:
scales_to_run = ["small", "medium", "full"]
results = {}
for scale_name in scales_to_run:
scale_def = SCALES[scale_name]
print(f"\n--- Benchmarking Scale: {scale_name} ({scale_def.approx_rows} rows) ---")
# Generate data for this scale
print(f"Generating synthetic pings for scale '{scale_name}'...")
subprocess.run([sys.executable, "generate_pings.py", "--scale", scale_name], check=True)
# Run CPU mode
print("Running CPU (pure pandas) baseline...")
cpu_timings = run_pipeline_subprocess(gpu_mode=False, scale_name=scale_name, temp_output_dir=temp_output_dir)
# Run GPU mode
print("Running GPU (cudf.pandas) pipeline...")
gpu_timings = run_pipeline_subprocess(gpu_mode=True, scale_name=scale_name, temp_output_dir=temp_output_dir)
results[scale_name] = {
"rows": scale_def.approx_rows,
"cpu": cpu_timings,
"gpu": gpu_timings
}
speedup = cpu_timings["total"] / max(gpu_timings["total"], 0.001)
print(f"Scale {scale_name} Results:")
print(f" CPU Total: {cpu_timings['total']:.3f}s")
print(f" GPU Total: {gpu_timings['total']:.3f}s")
print(f" Speedup: {speedup:.2f}x")
# Prepare full results object including metadata
import datetime
# Try fetching GPU name dynamically from nvidia-smi if available
gpu_hardware = "NVIDIA GPU (RAPIDS-accelerated)"
try:
gpu_info = subprocess.run(["nvidia-smi", "--query-gpu=name", "--format=csv,noheader"], capture_output=True, text=True)
if gpu_info.returncode == 0 and gpu_info.stdout.strip():
gpu_hardware = gpu_info.stdout.strip()
except Exception:
pass
full_results = {
"metadata": {
"timestamp": datetime.datetime.utcnow().strftime("%Y-%m-%d %H:%M:%S UTC"),
"hardware": gpu_hardware
}
}
# Copy all run metrics
for k, v in results.items():
full_results[k] = v
# Ensure results folder exists
results_dir = Path("results")
results_dir.mkdir(parents=True, exist_ok=True)
# Write JSON results to both data/ and results/
for out_path in [CFG.data_dir / "benchmark_results.json", results_dir / "benchmark_results.json"]:
with open(out_path, "w") as f:
json.dump(full_results, f, indent=2)
print(f"Saved benchmark results to {out_path}")
# Generate Markdown Table for README
md_table = []
md_table.append("| Scale | Data Rows | Stage | CPU Time (s) | GPU Time (s) | Speedup |")
md_table.append("|---|---|---|---|---|---|")
for scale_name, scale_res in results.items():
rows = scale_res["rows"]
for stage in ["load", "groupby_headway", "anomaly_scan", "scoring", "total"]:
cpu_t = scale_res["cpu"].get(stage, 0.0)
gpu_t = scale_res["gpu"].get(stage, 0.0)
speedup = cpu_t / max(gpu_t, 0.001)
stage_name = "**Total Insight Time**" if stage == "total" else stage
md_table.append(f"| {scale_name.capitalize()} | {rows} | {stage_name} | {cpu_t:.3f}s | {gpu_t:.3f}s | {speedup:.2f}x |")
print("\nBenchmark Markdown Table for README:")
print("\n".join(md_table))
print()
# Print the specific summary line required by specifications
headline_scale = scales_to_run[-1]
cpu_final = results[headline_scale]["cpu"]["total"]
gpu_final = results[headline_scale]["gpu"]["total"]
speedup_final = cpu_final / max(gpu_final, 0.001)
print("=================================================================")
print(f"Time-to-insight: {gpu_final:.2f}s (GPU) vs {cpu_final:.2f}s (CPU), {speedup_final:.2f}x speedup.")
print("=================================================================")
# Try plotting speedup chart to both data/ and results/
for out_chart in [CFG.data_dir / "speedup_chart.png", results_dir / "speedup_chart.png"]:
try:
plot_benchmark_results(results, out_chart)
except Exception as e:
print(f"Could not generate speedup chart at {out_chart}: {e}")
if __name__ == "__main__":
main()
|