tracinginsights's picture
download
raw
9.31 kB
#!/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.