From 8e8136d562f19e6e2f11f74cc58187a110eebbaf Mon Sep 17 00:00:00 2001 From: Sooyoung Cheong <64125280+c-sooyoung@users.noreply.github.com> Date: Fri, 14 Aug 2026 19:29:52 +0900 Subject: [PATCH] cleanup --- pipelines/batched_sobo.py | 146 +++++++------------------------------- 1 file changed, 25 insertions(+), 121 deletions(-) diff --git a/pipelines/batched_sobo.py b/pipelines/batched_sobo.py index 66c689d..289b3d3 100644 --- a/pipelines/batched_sobo.py +++ b/pipelines/batched_sobo.py @@ -1,47 +1,27 @@ import os import traceback import multiprocessing as mp -import shutil import bo import ptycho -def run_ptycho_worker( - worker_id, - gpu_token, - job_config, - metric, - run_id, - result_queue, -): - # This process, and MATLAB launched from it, - # can see exactly one GPU. +def run_ptycho_worker(worker_id, gpu_token, job_config, metric, run_id, result_queue): + + # This process, and MATLAB launched from it, can see exactly one GPU. os.environ["CUDA_VISIBLE_DEVICES"] = gpu_token try: - PTYCHOENGINE = ptycho.engines[job_config["ptycho"]["engine"]] - - ptycho_engine = PTYCHOENGINE(job_config) + ptycho_engine = ptycho.engines[job_config["ptycho"]["engine"]](job_config) ptycho_engine.run(run_id=run_id) y_value = ptycho_engine.metric(metric) - - result_queue.put( - (worker_id, y_value, None) - ) + result_queue.put((worker_id, y_value, None)) except Exception: - result_queue.put( - (worker_id, None, traceback.format_exc()) - ) + result_queue.put((worker_id, None, traceback.format_exc())) -def run_batch( - ctx, - gpu_tokens, - job_configs, - metric, - iteration, -): +def run_batch(ctx, gpu_tokens, job_configs, metric, iteration): + result_queue = ctx.Queue() processes = [] @@ -51,25 +31,12 @@ def run_batch( p = ctx.Process( target=run_ptycho_worker, - args=( - i, - gpu_tokens[i], - job_config, - metric, - run_id, - result_queue, - ), + args=(i, gpu_tokens[i], job_config, metric, run_id, result_queue), ) - p.start() processes.append(p) - # Four jobs are now running concurrently. - - results = [ - result_queue.get() - for _ in processes - ] + results = [result_queue.get() for _ in processes] # Synchronization barrier. for p in processes: @@ -80,111 +47,48 @@ def run_batch( for worker_id, _, error in results: if error is not None: - raise RuntimeError( - f"Ptycho worker {worker_id} failed:\n{error}" - ) + raise RuntimeError(f"Ptycho worker {worker_id} failed:\n{error}") return [y_value for _, y_value, _ in results] def sobo_pipeline(config): - result_dir = config["io"]["result_dir"] - if os.path.exists(result_dir): - shutil.rmtree(result_dir) - os.makedirs(result_dir, exist_ok=True) - RANDOM_ITERS = config["job"].get("random_iters", 0) SOBO_ITERS = config["job"].get("sobo_iters", 0) METRIC = config["bo"]["metric"] BO_BATCH = config["bo"]["batch"] - # SLURM should expose the four GPUs allocated to this job. visible = os.environ.get("CUDA_VISIBLE_DEVICES") - if visible is None: - raise RuntimeError( - "CUDA_VISIBLE_DEVICES is not set" - ) - - gpu_tokens = [ - token.strip() - for token in visible.split(",") - if token.strip() - ] - + raise RuntimeError("CUDA_VISIBLE_DEVICES is not set") + gpu_tokens = [token.strip() for token in visible.split(",") if token.strip()] if len(gpu_tokens) < BO_BATCH: - raise RuntimeError( - f"Expected {BO_BATCH} allocated GPUs, got {len(gpu_tokens)}" - ) + raise RuntimeError(f"Expected {BO_BATCH} allocated GPUs, got {len(gpu_tokens)}") # Explicitly use spawn for CUDA / MATLAB isolation. ctx = mp.get_context("spawn") - # --------------------------------------------------------- - # Random sampling - # --------------------------------------------------------- + ############################ RANDOM SAMPLING ############################### randombo = bo.RandomBOEngine(config) for j in range(RANDOM_ITERS): - print( - f"RANDOM sampling; iteration {j}", - flush=True, - ) - + print(f"RANDOM sampling; iteration {j}") job_configs = randombo.ask(n=BO_BATCH) + y_values = run_batch(ctx=ctx, gpu_tokens=gpu_tokens, job_configs=job_configs, metric=METRIC, iteration=j) + for job_config, y_value in zip(job_configs, y_values): + randombo.tell(job_config, y_value) - y_values = run_batch( - ctx=ctx, - gpu_tokens=gpu_tokens, - job_configs=job_configs, - metric=METRIC, - iteration=j, - ) - - # Only the parent touches BO state / train_x / train_y. - for job_config, y_value in zip( - job_configs, - y_values, - ): - randombo.tell( - job_config, - y_value, - ) - - # --------------------------------------------------------- - # SOBO - # --------------------------------------------------------- - + ############################ SOBO SAMPLING ############################### sobo = bo.SingleObjectiveBOEngine(config) - sobo.train_x = randombo.train_x sobo.train_y = randombo.train_y - for j in range(SOBO_ITERS): - iteration = RANDOM_ITERS + j - - print( - f"SOBO sampling; iteration {iteration}", - flush=True, - ) - + for j in range(RANDOM_ITERS, SOBO_ITERS+RANDOM_ITERS): + print(f"SOBO sampling iteration {j}") job_configs = sobo.ask(n=BO_BATCH) + y_values = run_batch(ctx=ctx, gpu_tokens=gpu_tokens, job_configs=job_configs, metric=METRIC, iteration=j) - y_values = run_batch( - ctx=ctx, - gpu_tokens=gpu_tokens, - job_configs=job_configs, - metric=METRIC, - iteration=iteration, - ) - - for job_config, y_value in zip( - job_configs, - y_values, - ): - sobo.tell( - job_config, - y_value, - ) + for job_config, y_value in zip(job_configs, y_values): + sobo.tell(job_config, y_value)