Buckets:
| #!/usr/bin/env python3 | |
| """Generate experiment variant scripts from the baseline R.py / Rp.py sources. | |
| Each generated script is a full, self-contained pipeline that: | |
| * writes all outputs under an env-defined EXP_ROOT (so a fresh dir forces | |
| full re-generation even when the canonical outputs already exist), and | |
| * skips the "wait for data" loop and goes straight to extraction. | |
| Run: python experiments/builder.py | |
| """ | |
| import os | |
| HERE = os.path.dirname(os.path.abspath(__file__)) | |
| REPO = os.path.dirname(HERE) | |
| BASE_DIR_LINE = ' base_dir = f"{event_name}/{session_name}"' | |
| BASE_DIR_PATCH = ''' _root = os.environ.get("EXP_ROOT", "").rstrip("/") | |
| base_dir = ( | |
| f"{_root}/{event_name}/{session_name}" | |
| if _root | |
| else f"{event_name}/{session_name}" | |
| ) | |
| os.makedirs(_root, exist_ok=True) if _root else None''' | |
| def make_baseline(src_name, dst_name): | |
| with open(os.path.join(REPO, src_name), "r", encoding="utf-8") as f: | |
| src = f.read() | |
| assert BASE_DIR_LINE in src, f"{src_name}: base_dir line not found" | |
| src = src.replace(BASE_DIR_LINE, BASE_DIR_PATCH, 1) | |
| # Replace the entire main() with a benchmark-friendly one (no wait loop). | |
| mark = "def main():" | |
| idx = src.index(mark) | |
| end_mark = 'if __name__ == "__main__":' | |
| end_idx = src.index(end_mark, idx) | |
| new_main = '''def main(): | |
| os.makedirs("cache", exist_ok=True) | |
| fastf1.Cache.enable_cache("cache") | |
| extractor = SeasonSessionExtractor(year=DEFAULT_YEAR) | |
| extractor.process_all() | |
| ''' | |
| src = src[:idx] + new_main + src[end_idx:] | |
| with open(os.path.join(HERE, dst_name), "w", encoding="utf-8") as f: | |
| f.write(src) | |
| print(f"wrote {dst_name}") | |
| def make_exp2_merge_once(): | |
| """Parallel pipeline but merge car+pos telemetry ONCE per driver, then | |
| slice per-lap (removes the repeated per-lap merge, the dominant cost).""" | |
| with open(os.path.join(HERE, "run_par.py"), "r", encoding="utf-8") as f: | |
| src = f.read() | |
| # A) Insert fast-path helpers after _lap_telemetry_or_none returns (before check_memory_usage). | |
| anchor = "def check_memory_usage(" | |
| helpers = '''def _lap_telemetry_from_merged(merged, selected): | |
| """Slice pre-merged telemetry for one lap (fast path, no re-merge).""" | |
| try: | |
| if merged is None or selected.empty: | |
| return None | |
| return merged.slice_by_lap(selected, interpolate_edges=True) | |
| except Exception: | |
| return None | |
| def _driver_merged_telemetry(driver_laps): | |
| """Build the interpolated merged car+pos telemetry once for a driver. | |
| Mirrors FastF1 ``Laps.get_telemetry`` merge but returns the merged object | |
| once so that per-lap slicing avoids repeating the expensive merge. | |
| Returns None when pos/car data cannot be merged. | |
| """ | |
| try: | |
| pos_data = driver_laps.get_pos_data(pad=1, pad_side="both") | |
| if pos_data is None or pos_data.empty or "Date" not in pos_data.columns: | |
| return None | |
| car_data = driver_laps.get_car_data(pad=1, pad_side="both") | |
| if ( | |
| car_data is None | |
| or car_data.empty | |
| or "Date" not in car_data.columns | |
| or len(car_data) < 3 | |
| ): | |
| return None | |
| drv_ahead = ( | |
| car_data.iloc[1:-1] | |
| .add_driver_ahead() | |
| .loc[:, ("DriverAhead", "DistanceToDriverAhead", "Date", "Time", "SessionTime")] | |
| ) | |
| car_data = car_data.add_distance().add_relative_distance() | |
| car_data = car_data.merge_channels(drv_ahead, frequency=None) | |
| return pos_data.merge_channels(car_data, frequency=None) | |
| except Exception: | |
| return None | |
| ''' | |
| src = src.replace(anchor, helpers + anchor, 1) | |
| # B) add merged param to _process_single_lap and use fast path. | |
| old_sig = " event_name: str,\n session_name: str,\n ) -> bool:" | |
| new_sig = " event_name: str,\n session_name: str,\n merged=None,\n ) -> bool:" | |
| assert old_sig in src, "sig not found" | |
| src = src.replace(old_sig, new_sig, 1) | |
| old_tel = " telemetry = _lap_telemetry_or_none(selected)" | |
| new_tel = (" telemetry = (\n" | |
| " _lap_telemetry_from_merged(merged, selected)\n" | |
| " if merged is not None\n" | |
| " else _lap_telemetry_or_none(selected)\n" | |
| " )") | |
| assert old_tel in src, "tel line not found" | |
| src = src.replace(old_tel, new_tel, 1) | |
| # C) process_driver: build merged once and pass into each lap. | |
| loop_old = ''' existing = ( | |
| set(os.listdir(driver_dir)) | |
| if os.path.isdir(driver_dir) | |
| else set() | |
| ) | |
| for lap_number in lap_numbers: | |
| fname = f"{lap_number}_tel.json" | |
| if fname in existing: | |
| continue | |
| self._process_single_lap( | |
| driver, lap_number, driver_dir, driver_laps, event_name, session_name | |
| )''' | |
| loop_new = ''' existing = ( | |
| set(os.listdir(driver_dir)) | |
| if os.path.isdir(driver_dir) | |
| else set() | |
| ) | |
| merged = _driver_merged_telemetry(driver_laps) | |
| for lap_number in lap_numbers: | |
| fname = f"{lap_number}_tel.json" | |
| if fname in existing: | |
| continue | |
| self._process_single_lap( | |
| driver, lap_number, driver_dir, driver_laps, | |
| event_name, session_name, merged | |
| )''' | |
| assert loop_old in src, "loop not found" | |
| src = src.replace(loop_old, loop_new, 1) | |
| with open(os.path.join(HERE, "exp2_merge_once.py"), "w", encoding="utf-8") as f: | |
| f.write(src) | |
| print("wrote exp2_merge_once.py") | |
| def make_exp3_workers(): | |
| """exp2 + raise the default worker cap so all 22 drivers run in one wave.""" | |
| with open(os.path.join(HERE, "exp2_merge_once.py"), "r", encoding="utf-8") as f: | |
| src = f.read() | |
| old = " else max(1, min(16, (os.cpu_count() or 2)))" | |
| new = " else max(1, min(32, (os.cpu_count() or 2)))" | |
| assert old in src, "worker cap line not found" | |
| src = src.replace(old, new, 1) | |
| with open(os.path.join(HERE, "exp3_workers.py"), "w", encoding="utf-8") as f: | |
| f.write(src) | |
| print("wrote exp3_workers.py") | |
| def make_exp5_saturate(): | |
| """Byte-identical: keep the exact per-lap FastF1 get_telemetry pipeline (as | |
| Rp.py) but raise the default worker cap so all drivers run in a single | |
| wave. This is the fastest variant that still produces byte-identical | |
| output to R.py.""" | |
| with open(os.path.join(HERE, "run_par.py"), "r", encoding="utf-8") as f: | |
| src = f.read() | |
| old = " else max(1, min(16, (os.cpu_count() or 2)))" | |
| new = " else max(1, min(32, (os.cpu_count() or 2)))" | |
| assert old in src, "worker cap line not found" | |
| src = src.replace(old, new, 1) | |
| with open(os.path.join(HERE, "exp5_saturate.py"), "w", encoding="utf-8") as f: | |
| f.write(src) | |
| print("wrote exp5_saturate.py") | |
| def make_exp4_fastslice(): | |
| """exp3 + replace the per-lap interpolated slice (which re-merges edges) | |
| with a vectorised boolean-mask slice over the once-built merged telemetry.""" | |
| with open(os.path.join(HERE, "exp3_workers.py"), "r", encoding="utf-8") as f: | |
| src = f.read() | |
| old = '''def _lap_telemetry_from_merged(merged, selected): | |
| """Slice pre-merged telemetry for one lap (fast path, no re-merge).""" | |
| try: | |
| if merged is None or selected.empty: | |
| return None | |
| return merged.slice_by_lap(selected, interpolate_edges=True) | |
| except Exception: | |
| return None''' | |
| new = '''def _lap_telemetry_from_merged(merged, selected): | |
| """Slice pre-merged telemetry for one lap using a boolean mask. | |
| Avoids the per-lap ``slice_by_lap(interpolate_edges=True)`` call, which | |
| re-merges interpolated edge samples for every lap. Drops the two | |
| interpolated boundary samples per lap (row-1 at each edge) but keeps all | |
| real channel values, so the produced telemetry is equivalent and file | |
| counts are unchanged - while being dramatically faster. | |
| """ | |
| try: | |
| if merged is None or selected.empty: | |
| return None | |
| st = selected["LapStartTime"].iloc[0] | |
| et = selected["Time"].iloc[0] | |
| t = merged["SessionTime"].to_numpy() | |
| idx = np.where((t >= st) & (t <= et))[0] | |
| if len(idx) == 0: | |
| return None | |
| data = merged.iloc[idx].copy() | |
| if "Time" in data.columns: | |
| data.loc[:, "Time"] = merged["SessionTime"].iloc[idx] - st | |
| return data | |
| except Exception: | |
| return None''' | |
| assert old in src, "fast slice fn anchor not found" | |
| src = src.replace(old, new, 1) | |
| with open(os.path.join(HERE, "exp4_fastslice.py"), "w", encoding="utf-8") as f: | |
| f.write(src) | |
| print("wrote exp4_fastslice.py") | |
| if __name__ == "__main__": | |
| make_baseline("R.py", "run_seq.py") # sequential baseline | |
| make_baseline("Rp.py", "run_par.py") # parallel baseline | |
| make_exp2_merge_once() | |
| make_exp3_workers() | |
| make_exp4_fastslice() | |
| make_exp5_saturate() | |
Xet Storage Details
- Size:
- 9.31 kB
- Xet hash:
- 38723fa0aac8edff22d88486f28bf3e9fc12c0e306e7e454cf0d6af90812eeda
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.