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()