Fabio Geraci
12/01/2025, 9:17 AM@task
def build_tmp_job_config
@task
def load_raw_dataset
@task
def transform_dataset
@task
def load_dataset_from_tag
@task
def submit_job
@flow
def hpc_training_flow():
# Use JOB_CONFIG_PATH as template; write /tmp/job_config.yml locally
tmp_job_config_path = build_tmp_job_config(JOB_CONFIG_PATH) # /tmp/job_config.yml
load_raw_dataset or load_dataset_from_tag
if load_raw_dataset:
transform_dataset
# 2) Submit job
job_id = submit_job(tmp_job_config_path)
final_status = monitor_job(job_id)
return final_status
if __name__ == "__main__":
hpc_training_flow()Marvin
12/01/2025, 9:17 AMMarvin
12/01/2025, 9:19 AMfrom prefect import flow, task
@task
def build_tmp_job_config(path: str) -> str:
# write to /tmp and return the tmp path
return "/tmp/job_config.yml"
@task
def load_raw_dataset() -> str:
return "raw_dataset"
@task
def transform_dataset(ds: str) -> str:
return f"transformed({ds})"
@task
def load_dataset_from_tag(tag: str) -> str:
return f"dataset_from_tag({tag})"
@task
def submit_job(job_config_path: str, dataset: str) -> str:
return "job-123"
@task
def monitor_job(job_id: str) -> str:
return "succeeded"
@flow
def hpc_training_flow(
job_config_path: str,
load_from_tag: bool = False,
tag: str | None = None,
):
tmp_job_config_path = build_tmp_job_config(job_config_path)
if load_from_tag and tag:
dataset = load_dataset_from_tag(tag)
else:
raw = load_raw_dataset()
dataset = transform_dataset(raw)
job_id = submit_job(tmp_job_config_path, dataset)
final_status = monitor_job(job_id)
return final_status
if __name__ == "__main__":
hpc_training_flow(job_config_path="/path/to/JOB_CONFIG_PATH", load_from_tag=False)
Decide based on a task result at runtime
from prefect import flow, task
@task
def should_use_tag(tag: str | None) -> bool:
# e.g., check tag exists in a registry
return bool(tag)
@flow
def hpc_training_flow(job_config_path: str, tag: str | None = None):
tmp_job_config_path = build_tmp_job_config(job_config_path)
use_tag = should_use_tag(tag)
if use_tag:
dataset = load_dataset_from_tag(tag) # type: ignore[arg-type]
else:
raw = load_raw_dataset()
dataset = transform_dataset(raw)
job_id = submit_job(tmp_job_config_path, dataset)
return monitor_job(job_id)
Notes
- Calling a task like result = task_fn(args) runs it and returns its result. Use .submit() only if you want to run tasks concurrently; if you branch on a submitted task’s output, call .result() to get the boolean for the if.
- No special branching operators are needed—just if/elif/else.