quinnlue's picture
MAE cross-objective mixture-ranking campaign
ffdcfe7 verified
Raw
History Blame Contribute Delete
4.03 kB
"""GPU-queue launcher for eatmap pretraining jobs.
One process per GPU pulling from a shared job list, because every campaign here
is embarrassingly parallel across mixtures and the box has 4 devices. Jobs are
skipped if their output directory already carries a finished run.json, so a
crashed launcher can be re-run without redoing work -- which matters on an
instance whose filesystem does not survive a recycle.
Usage:
python launch.py jobs.json [--gpus 0,1,2,3] [--dry-run]
jobs.json is a list of {name, configs[], sets[]} records.
"""
import argparse
import json
import queue
import subprocess
import threading
import time
from pathlib import Path
REPO = "/workspace/code/eat-map-regmix"
RUNS = Path("/workspace/runs")
LOGS = Path("/workspace/logs")
def finished(out: Path) -> bool:
status = out / "status.json"
if not status.exists():
return False
try:
# runner.py writes {"completed": true, "final_step": ...} on a clean exit.
return bool(json.loads(status.read_text()).get("completed"))
except Exception:
return False
def worker(gpu: int, jobs: queue.Queue, results: list, lock: threading.Lock, dry: bool):
while True:
try:
job = jobs.get_nowait()
except queue.Empty:
return
name = job["name"]
out = RUNS / name
if finished(out):
with lock:
print(f"[gpu{gpu}] skip {name} (already finished)")
results.append((name, "skipped", 0.0))
jobs.task_done()
continue
cmd = ["python", "-m", "eatmap.cli.pretrain", "--output-dir", str(out)]
for c in job["configs"]:
cmd += ["--config", c]
for s in job.get("sets", []):
cmd += ["--set", s]
# Job names carry a campaign prefix ("mae-lr/p16-lr0.003"), so the log
# path is nested and its parent has to exist before the open() below.
log = LOGS / f"{name}.log"
log.parent.mkdir(parents=True, exist_ok=True)
with lock:
print(f"[gpu{gpu}] start {name}")
if dry:
with lock:
print(" " + " ".join(cmd))
jobs.task_done()
continue
t0 = time.time()
with open(log, "w") as fh:
proc = subprocess.run(
cmd, cwd=REPO, stdout=fh, stderr=subprocess.STDOUT,
env={**__import__("os").environ, "CUDA_VISIBLE_DEVICES": str(gpu),
"PYTHONPATH": REPO},
)
dt = time.time() - t0
state = "ok" if proc.returncode == 0 else f"FAIL({proc.returncode})"
with lock:
print(f"[gpu{gpu}] {state} {name} in {dt/60:.1f} min -> {log}")
results.append((name, state, dt))
jobs.task_done()
def main():
ap = argparse.ArgumentParser()
ap.add_argument("jobs")
ap.add_argument("--gpus", default="0,1,2,3")
ap.add_argument("--dry-run", action="store_true")
args = ap.parse_args()
spec = json.loads(Path(args.jobs).read_text())
jobs: queue.Queue = queue.Queue()
for job in spec:
jobs.put(job)
gpus = [int(g) for g in args.gpus.split(",")]
print(f"{len(spec)} jobs over {len(gpus)} gpus")
results: list = []
lock = threading.Lock()
threads = [threading.Thread(target=worker, args=(g, jobs, results, lock, args.dry_run))
for g in gpus]
t0 = time.time()
for t in threads:
t.start()
for t in threads:
t.join()
ok = [r for r in results if r[1] == "ok"]
bad = [r for r in results if r[1].startswith("FAIL")]
print(f"\nwall {(time.time()-t0)/60:.1f} min | ok {len(ok)} | skipped "
f"{len([r for r in results if r[1]=='skipped'])} | failed {len(bad)}")
for name, state, _ in bad:
print(f" {state} {name} -> /workspace/logs/{name}.log")
if ok:
mean = sum(r[2] for r in ok) / len(ok)
print(f"mean run time {mean/60:.1f} min")
if __name__ == "__main__":
main()