← FlyteCONTENT HISTORY

Update to Flyte

Snapshot Sep 30, 2026 · 22:59 UTC · version 1.0.1

Collection source: not recorded for this historical snapshot.

WHAT CHANGED · RULE-BASED ANALYSIS

First saved snapshot

No earlier snapshot is available to establish a change.

Compare saved observations

Download comparison JSON
Full technical diff · 0 changed fields
Full snapshot data
{
  "description": "Migrates Slurm (sbatch/srun) HPC workloads to Flyte 2 (the flyte Python SDK) — job scripts become typed tasks, `#SBATCH` pragmas become TaskEnvironment config, job arrays become flyte.map, and multi-node training becomes a clustered task environment. Use when porting an HPC or supercomputer cluster workload off Slurm, translating `#SBATCH` pragmas, or replacing sbatch chains, job arrays, and module load with Flyte. Trigger words are sbatch, srun, SLURM, `#SBATCH`, HPC, job array, partition, squeue, sinfo, module load, scancel, supercomputer, cluster migration.",
  "included_files": [],
  "name": "flyte-migrate-slurm",
  "skill_md_contents": "---\nname: flyte-migrate-slurm\ndescription: Migrates Slurm (sbatch/srun) HPC workloads to Flyte 2 (the flyte Python SDK) — job scripts become typed tasks, `#SBATCH` pragmas become TaskEnvironment config, job arrays become flyte.map, and multi-node training becomes a clustered task environment. Use when porting an HPC or supercomputer cluster workload off Slurm, translating `#SBATCH` pragmas, or replacing sbatch chains, job arrays, and module load with Flyte. Trigger words are sbatch, srun, SLURM, `#SBATCH`, HPC, job array, partition, squeue, sinfo, module load, scancel, supercomputer, cluster migration.\n---\n\n# Slurm to Flyte 2 Migration Skill\n\nSlurm schedules jobs. Somewhere along the way ML work stopped being jobs and became pipelines — a data prep step, a training step, an eval step, a sweep over configs, each with different hardware and failure characteristics. A Slurm job is a bash script with `#SBATCH` pragmas at the top; in Flyte 2 that script becomes a typed Python function decorated with `@env.task`, a pipeline is just a task that calls other tasks, and control flow is plain Python. This skill maps the Slurm surface area onto the Flyte 2 SDK and is honest about the places where Slurm still wins.\n\n## Grounding References\n\n| Resource          | URL                                                              |\n| ----------------- | ---------------------------------------------------------------- |\n| Migration guide   | https://www.union.ai/docs/v2/flyte/user-guide/migration/flyte-2/ |\n| Official docs     | https://www.union.ai/docs/v2/flyte                               |\n| Docs index (LLMs) | https://www.union.ai/docs/v2/flyte/llms.txt                      |\n| SDK API reference | https://www.union.ai/docs/v2/union/api-reference/flyte-sdk/      |\n| Example code      | https://github.com/unionai/unionai-examples                      |\n| Flyte MCP tools   | Available via `flyte-mcp` server                                 |\n\n## Three mental shifts\n\nAlmost every point of friction for a Slurm migrant traces back to one of these three.\n\n### 1. Modules become images\n\nOn Slurm you `module load cuda/12.1`, activate a venv on NFS, and hope the login node and the compute node agree. In Flyte 2 the environment is declared once, in Python, as a `flyte.Image`. There is no Dockerfile to write: images are built remotely and content-hashed, so an unchanged spec is a cache hit and a changed one rebuilds automatically.\n\n```python\nimage = (\n    flyte.Image.from_debian_base(python_version=(3, 12))\n    .with_pip_packages(\"torch\", \"transformers\", \"datasets\")\n    .with_env_vars({\"HF_HUB_ENABLE_HF_TRANSFER\": \"1\"})\n)\n```\n\n### 2. The shared filesystem becomes explicit data\n\nThere is no shared `/home` or `/scratch` that every node sees. Instead you pass `flyte.io.File` and `flyte.io.Dir` — typed references to object storage that stream rather than copy. This is more typing than `/scratch/$USER/run17/ckpt.pt`, and in exchange you get lineage for free: every input and output of every task is recorded, so \"which dataset produced this checkpoint\" is a question the system answers instead of a question you grep for. If you genuinely need a parallel filesystem (Lustre, GPFS, FSx), mount it through a CSI driver and a `flyte.PodTemplate` — that path stays open.\n\n### 3. Job scripts become functions\n\n`sbatch` chains held together by `--dependency=afterok`, sentinel files on NFS, and a cron job that polls for them all collapse into ordinary Python: call a function, await it, branch on the result. Binaries you can't or won't rewrite in Python still run — as container tasks with typed inputs and outputs.\n\n## Migration Cheat Sheet\n\n| Slurm                                             | Flyte 2                                                     |\n| ------------------------------------------------- | ----------------------------------------------------------- |\n| `sbatch train.sh`                                 | `flyte run train.py main`                                   |\n| `srun --pty python train.py` (interactive)        | `flyte run --local train.py main`, or a devbox              |\n| `#SBATCH --gres=gpu:a100:8`                       | `flyte.Resources(gpu=\"A100:8\")`                             |\n| `#SBATCH --cpus-per-task=16 --mem=64G`            | `flyte.Resources(cpu=16, memory=\"64Gi\")`                    |\n| `#SBATCH --tmp=100G`                              | `flyte.Resources(disk=\"100Gi\")`                             |\n| `#SBATCH --array=0-999`                           | `flyte.map(step, range(1000), concurrency=200)`             |\n| `#SBATCH --nodes=4 --ntasks-per-node=8`           | `ClusteredTaskEnvironment(replicas=4, nproc_per_node=8)`    |\n| `#SBATCH --partition=gpu --qos=high`              | `queue=\"gpu-high\"` (queue must exist in cluster config)     |\n| `#SBATCH --requeue`                               | `retries=3` plus `interruptible=True` for spot              |\n| `#SBATCH --time=04:00:00`                         | `timeout=timedelta(hours=4)`                                |\n| `#SBATCH --begin=...` / a crontab entry           | `flyte.Trigger(...)` with `flyte.Cron(...)`                 |\n| `#SBATCH --dependency=afterok:$JOBID`             | Plain Python — call the next task after the first returns   |\n| `module load cuda && source venv/bin/activate`    | `flyte.Image.from_debian_base().with_pip_packages(...)`     |\n| `$SLURM_PROCID`, `$SLURM_NNODES`, `$SLURM_NTASKS` | `flyte.ctx().rank`, `.nnodes`, `.world_size`                |\n| `$SLURM_ARRAY_TASK_ID`                            | The argument you mapped over                                |\n| `/scratch/$USER/data.parquet`                     | `flyte.io.File` / `flyte.io.Dir` passed between tasks       |\n| `squeue`, `sacct`                                 | `flyte get run`, `flyte get logs`, the UI                   |\n| `scancel <jobid>`                                 | `flyte stop`, or abort from the UI                          |\n| `srun --pty bash` / `ssh node042`                 | `flyte run --debug ...`, or `flyte debug <run-name>` (beta) |\n\n## The job script becomes a task\n\nThis is the canonical translation. Everything above the `module load` line is configuration and moves onto the `TaskEnvironment` or the task decorator; everything below it is the function body.\n\n### Slurm\n\n```bash\n#!/bin/bash\n#SBATCH --job-name=train\n#SBATCH --partition=gpu\n#SBATCH --gres=gpu:a100:8\n#SBATCH --cpus-per-task=16\n#SBATCH --mem=64G\n#SBATCH --time=04:00:00\n#SBATCH --requeue\n\nmodule load cuda/12.1\nsource ~/venvs/train/bin/activate\nsrun python train.py --lr 3e-4\n```\n\n### Flyte 2\n\n```python\nfrom datetime import timedelta\n\nimport flyte\nfrom flyte.io import File\n\nenv = flyte.TaskEnvironment(\n    name=\"training\",\n    image=flyte.Image.from_debian_base(python_version=(3, 12)).with_pip_packages(\"torch\"),\n    resources=flyte.Resources(cpu=16, memory=\"64Gi\", gpu=\"A100:8\"),\n)\n\n\n@env.task(retries=3, timeout=timedelta(hours=4))\nasync def train(lr: float = 3e-4) -> File:\n    import torch\n\n    model = build_model().cuda()\n    optimizer = torch.optim.AdamW(model.parameters(), lr=lr)\n    ...  # the same training loop that was in train.py\n    torch.save(model.state_dict(), \"model.pt\")\n    return await File.from_local(\"model.pt\")\n```\n\n```bash\nflyte run train.py train --lr 3e-4     # remote (the sbatch equivalent)\nflyte run --local train.py train       # same code, on your laptop\nflyte deploy train.py env              # register the environment + triggers\n```\n\nThe `#SBATCH --time` line has a richer counterpart than a single number. `timeout=` accepts a `timedelta`, an int number of seconds, or a `flyte.Timeout` object that separates the budgets Slurm collapses into one:\n\n```python\n@env.task(\n    retries=2,\n    timeout=flyte.Timeout(\n        max_runtime=timedelta(hours=4),      # per-attempt wall clock\n        max_queued_time=timedelta(minutes=30),  # fail fast if capacity never appears\n        deadline=timedelta(hours=10),        # absolute budget across all attempts\n    ),\n)\nasync def train(...) -> File: ...\n```\n\n## Job arrays and sweeps\n\nA Slurm array plus a results-collection script is two artifacts held together by a filename convention. In Flyte 2 the fan-out and the reduction live in the same function.\n\n### Slurm\n\n```bash\n#!/bin/bash\n#SBATCH --array=0-63\n#SBATCH --gres=gpu:1\nCONFIG=$(sed -n \"$((SLURM_ARRAY_TASK_ID + 1))p\" configs.txt)\npython train.py --config \"$CONFIG\" --out \"/scratch/$USER/sweep/$SLURM_ARRAY_TASK_ID.json\"\n# ...then a separate job, after the array drains, to read 64 JSON files and pick a winner.\n```\n\n### Flyte 2\n\n```python\nimport asyncio\n\nimport flyte\n\nenv = flyte.TaskEnvironment(\n    name=\"sweep\",\n    image=flyte.Image.from_debian_base().with_pip_packages(\"torch\"),\n    resources=flyte.Resources(cpu=8, memory=\"32Gi\", gpu=\"L4:1\"),\n)\n\n\n@env.task\nasync def train_one(lr: float, batch_size: int) -> float:\n    ...  # returns a validation metric\n    return val_loss\n\n\n@env.task\nasync def sweep() -> dict:\n    configs = [(lr, bs) for lr in (1e-4, 3e-4, 1e-3) for bs in (16, 32, 64)]\n    # Small fan-out: gather is the most direct translation of an array job.\n    losses = await asyncio.gather(*[train_one(lr, bs) for lr, bs in configs])\n    # Picking the winner is plain Python — no second job, no sentinel files.\n    best = min(range(len(losses)), key=lambda i: losses[i])\n    return {\"lr\": configs[best][0], \"batch_size\": configs[best][1], \"loss\": losses[best]}\n```\n\nFor a 10,000-element array, `flyte.map` gives you a bounded fan-out with a concurrency cap — the `%` throttle in `--array=0-9999%200`:\n\n```python\n@env.task\nasync def big_sweep(n: int = 10_000) -> float:\n    # flyte.map returns a generator; wrap it in list() to materialize.\n    losses = list(flyte.map(train_one_indexed, range(n), concurrency=200))\n    return min(losses)\n```\n\n## Job dependencies become function calls\n\n`--dependency=afterok:$JOBID` is the construct that most often turns into a shell script full of `sbatch --parsable` and `awk`. It has no counterpart in Flyte 2 because it doesn't need one: awaiting a task _is_ the dependency, and the value it returns _is_ the handoff.\n\n### Slurm\n\n```bash\nPREP=$(sbatch --parsable prep.sh)\nTRAIN=$(sbatch --parsable --dependency=afterok:$PREP train.sh)\nsbatch --dependency=afterok:$TRAIN eval.sh\n```\n\n### Flyte 2\n\n```python\nfrom flyte.io import Dir, File\n\n\n@env.task\nasync def pipeline(raw: Dir) -> float:\n    prepared = await prep(raw)         # runs first\n    model = await train(prepared)      # waits for prep, gets its output directly\n    return await evaluate(model, prepared)\n```\n\nBranching that Slurm can't express at all — `afterok` on one job but `afternotok` on another, or \"promote only if the eval didn't regress\" — is just `try`/`except` and `if`:\n\n```python\n@env.task\nasync def pipeline_with_gate(raw: Dir) -> str:\n    prepared = await prep(raw)\n    try:\n        model = await train(prepared)\n    except Exception:\n        model = await train_with_fallback_config(prepared)\n    if await evaluate(model, prepared) < 0.85:\n        return \"held back\"\n    await promote(model)\n    return \"promoted\"\n```\n\n## Data: from `/scratch` to typed references\n\nThe Slurm version writes to a path both jobs happen to agree on. The Flyte version passes a value.\n\n### Slurm\n\n```bash\n# prep.sh\npython prep.py --in /scratch/$USER/raw --out /scratch/$USER/prepared\n\n# train.sh — coupled to prep.sh only by this string\npython train.py --data /scratch/$USER/prepared --ckpt /scratch/$USER/ckpt.pt\n```\n\n### Flyte 2\n\n```python\nimport flyte\nfrom flyte.io import Dir, File\n\n\n@env.task\nasync def prep(raw: Dir) -> Dir:\n    local = await raw.download()\n    ...  # write outputs into ./prepared\n    return await Dir.from_local(\"prepared\")\n\n\n@env.task\nasync def train(prepared: Dir) -> File:\n    local = await prepared.download()  # streams from object storage to this pod\n    ...\n    return await File.from_local(\"ckpt.pt\")\n```\n\nScratch space _inside_ a task is still just the local filesystem — request it with `flyte.Resources(disk=\"100Gi\")` and use `os.getcwd()`. What changes is that anything another task needs must leave as a typed output. If a parallel filesystem is non-negotiable (a dataset too large or too latency-sensitive to stream), mount it with a CSI driver through `pod_template=flyte.PodTemplate.from_spec(pod_spec_with_lustre_pvc)` on the `TaskEnvironment`.\n\n## Multi-node training\n\n`--nodes=4 --ntasks-per-node=8` with `srun` as the launcher becomes a `ClusteredTaskEnvironment`. It launches its replicas as a single Kubernetes JobSet, runs `torchrun` rendezvous across them, and exposes the standard `RANK` / `WORLD_SIZE` / `MASTER_ADDR` environment variables — so training code written for `torchrun` needs no changes. The same values are on `flyte.ctx()`.\n\n### Slurm\n\n```bash\n#!/bin/bash\n#SBATCH --nodes=4\n#SBATCH --ntasks-per-node=8\n#SBATCH --gres=gpu:h100:8\nexport MASTER_ADDR=$(scontrol show hostnames \"$SLURM_JOB_NODELIST\" | head -n1)\nsrun python -m torch.distributed.run --nnodes=4 --nproc_per_node=8 pretrain.py\n```\n\n### Flyte 2\n\n```python\nimport flyte\nfrom flyte.clustered import ClusteredTaskEnvironment, ClusterFailurePolicy, TorchRun\n\nenv = ClusteredTaskEnvironment(\n    name=\"pretrain\",\n    image=image,\n    resources=flyte.Resources(cpu=16, memory=\"64Gi\", gpu=\"H100:8\", shm=\"auto\"),\n    replicas=4,          # pods == nodes\n    nproc_per_node=8,    # processes per pod  =>  world size 32\n    runtime=TorchRun(rdzv_backend=\"c10d\"),  # \"static\" relies on JobSet restarts instead\n    failure_policy=ClusterFailurePolicy(max_restarts=2, restart_on_host_maintenance=True),\n)\n\n\n@env.task\nasync def pretrain(steps: int = 10_000) -> File:\n    import torch\n    import torch.distributed as dist\n\n    ctx = flyte.ctx()\n    torch.cuda.set_device(ctx.local_rank or 0)\n    dist.init_process_group(backend=\"nccl\")  # torchrun already set RANK/WORLD_SIZE/MASTER_ADDR\n    print(f\"rank {ctx.rank}/{ctx.world_size} on node {ctx.node_rank}/{ctx.nnodes}\", flush=True)\n\n    ...  # the same DDP/FSDP loop you ran under srun\n\n    dist.barrier()\n    dist.destroy_process_group()\n    # Only rank 0 has anything to return.\n    return await File.from_local(\"ckpt.pt\")\n```\n\n`restart_on_host_maintenance=True` is the piece with no Slurm analogue: a node reclaimed by the cloud provider (spot reclaim, host maintenance, drain) restarts the whole set _for free_, leaving the `max_restarts` budget untouched, so an unreliable cluster can't burn the budget you reserved for real bugs.\n\nA clustered task is a worker, not a driver — it cannot launch subtasks. To compose distributed steps, orchestrate from a plain `TaskEnvironment` that declares `depends_on=[clustered_env]`:\n\n```python\ndriver_env = flyte.TaskEnvironment(\n    name=\"driver\",\n    image=image,\n    resources=flyte.Resources(cpu=1, memory=\"1Gi\"),\n    depends_on=[env],  # without this, awaiting pretrain() fails on image-cache lookup\n)\n\n\n@driver_env.task\nasync def main(steps: int = 10_000) -> float:\n    ckpt = await pretrain(steps)   # JobSet #1\n    return await evaluate(ckpt)    # JobSet #2\n```\n\nEphemeral, per-task Ray / Spark / Dask clusters (via the `flyteplugins-ray`, `flyteplugins-spark`, and `flyteplugins-dask` integrations) replace the long-lived clusters an HPC site usually stands up by hand — they exist for the duration of the task and are torn down with it.\n\n## Fault tolerance: `--requeue`, decomposed\n\n`#SBATCH --requeue` restarts the job from the top and hopes you wrote your own checkpoint logic. Flyte 2 splits that into four independent mechanisms you compose.\n\n**Retries, declaratively.** `retries=3` on the decorator, or a `RetryStrategy` when you want backoff:\n\n```python\n@env.task(\n    retries=flyte.RetryStrategy(\n        count=4,\n        backoff=flyte.Backoff(base=timedelta(seconds=10), factor=2.0, cap=timedelta(minutes=5)),\n    ),\n)\nasync def flaky_download(url: str) -> File: ...\n```\n\n**Spot capacity, safely.** `interruptible=True` says the task may run on preemptible instances. Preemptions are tracked as _system_ failures and do not consume your retry budget, and the final attempt falls back to on-demand — so `interruptible=True, retries=2` means two spot attempts and one guaranteed on-demand attempt. Set it on the environment, override it per task:\n\n```python\nenv = flyte.TaskEnvironment(name=\"sweep\", image=image, interruptible=True)\n\n\n@env.task(interruptible=False)  # the one step you don't want preempted\nasync def publish(model: File) -> str: ...\n```\n\n**Checkpoints that survive node changes.** Checkpoints go to object storage, not `/scratch`, so a retry resumes on whatever node it lands on. No shared filesystem required:\n\n```python\n@env.task(retries=3)\nasync def train(n_epochs: int = 100) -> int:\n    checkpoint = flyte.ctx().checkpoint\n    path = await checkpoint.load()            # None on the first attempt\n    start = int(path.read_bytes()) if path else 0\n\n    for epoch in range(start, n_epochs):\n        ...\n        await checkpoint.save(f\"{epoch + 1}\".encode())\n    return n_epochs\n```\n\n**Durable function calls.** `@flyte.trace` records the result of an individual function call inside a task. On a retry, recorded calls are skipped instead of re-executed — the granularity Slurm has no way to express:\n\n```python\n@flyte.trace\nasync def call_expensive_api(prompt: str) -> str:\n    ...  # on retry, an already-recorded call replays instead of re-running\n```\n\n**Caching.** Task-level caching is keyed on the code and the inputs, so rerunning a twelve-hour pipeline after fixing step nine starts at step nine:\n\n```python\nenv = flyte.TaskEnvironment(name=\"etl\", image=image, cache=\"auto\")\n```\n\n## Warm pools\n\nSlurm feels fast at submit time because the allocation is already running — `srun` inside an existing allocation starts in milliseconds. `flyte.ReusePolicy` is the equivalent: containers stay warm between tasks, keeping in-memory state, so a model loaded once serves thousands of task invocations.\n\n```python\nenv = flyte.TaskEnvironment(\n    name=\"scorer\",\n    image=image,\n    resources=flyte.Resources(cpu=4, memory=\"16Gi\", gpu=\"L4:1\"),\n    reusable=flyte.ReusePolicy(\n        replicas=(2, 10),   # autoscaling range; a bare int pins the count\n        concurrency=4,      # concurrent tasks per replica (async tasks only)\n        idle_ttl=300,       # shut the environment down after 5 idle minutes\n    ),\n)\n```\n\nThe caveat is the same one that bites long-lived Slurm allocations: the process outlives the task. Treat module-level globals and caches deliberately, and don't let one task's state leak into the next.\n\n## Existing binaries\n\nA Fortran solver, a C++ simulator, a genomics tool — anything you're not rewriting runs as a `ContainerTask` with typed inputs and outputs. Inputs are staged into `input_data_dir`, and whatever the command writes into `output_data_dir` is read back as the declared types.\n\n```python\nfrom flyte.extras import ContainerTask\nfrom flyte.io import File\n\nalign = ContainerTask(\n    name=\"align_reads\",\n    image=\"quay.io/biocontainers/bwa:0.7.17\",\n    resources=flyte.Resources(cpu=16, memory=\"64Gi\"),\n    inputs={\"reference\": File, \"reads\": File},\n    outputs={\"alignment\": File},\n    input_data_dir=\"/var/inputs\",\n    output_data_dir=\"/var/outputs\",\n    file_input_layout=\"NAMED_DIR\",  # preserves original filenames + extensions\n    command=[\n        \"/bin/sh\", \"-c\",\n        \"bwa mem /var/inputs/reference/* /var/inputs/reads/* > /var/outputs/alignment\",\n    ],\n)\n\n\n@env.task\nasync def main(reference: File, reads: File) -> File:\n    return await align(reference=reference, reads=reads)\n```\n\n## Scheduling, monitoring, and interactive work\n\n`#SBATCH --begin=` and the crontab that wraps most recurring HPC work become a `flyte.Trigger` attached to the task and deployed with it:\n\n```python\nfrom datetime import datetime\n\nnightly = flyte.Trigger(\n    name=\"nightly_retrain\",\n    automation=flyte.Cron(\"0 2 * * *\", timezone=\"America/Los_Angeles\"),\n    inputs={\"start_time\": flyte.TriggerTime, \"lr\": 3e-4},\n    auto_activate=True,\n)\n\n\n@env.task(triggers=nightly)\nasync def retrain(start_time: datetime, lr: float) -> File: ...\n```\n\nFor monitoring, `squeue` and `sacct` become `flyte get run` and `flyte get logs` (or the UI, which shows the pipeline structure rather than a flat job list):\n\n```bash\nflyte get run                                  # like squeue\nflyte get run <run_name>                       # detail for one run\nflyte get logs <run_name> --attempt 0          # like sacct + tailing a slurm-*.out\nflyte stop <run_name>                          # like scancel\n```\n\n`srun --pty bash` and `ssh node042` have two counterparts. `flyte run --debug` opens a browser-based VS Code session in the task pod:\n\n```bash\nflyte run --debug train.py train --lr 3e-4\n```\n\n```python\nrun = flyte.with_runcontext(debug=True).run(train, lr=3e-4)\nprint(run.get_debug_url())\n```\n\nSSH into the running task is available in beta (requires `flyteplugins-union`) and is closer to the muscle memory of `ssh` onto a compute node:\n\n```bash\nflyte debug <run-name> --write-config\nssh flyte-debug\n```\n\n## Migration order that works\n\nDo not start with the thing Slurm does best.\n\n1. **Pipeline-shaped work first** — data processing, evals, sweeps, batch inference, RL rollouts. These are multi-step, embarrassingly parallel, and failure-prone in boring ways, so composition, caching, and retries pay off on day one. They also exercise images and data plumbing on workloads where a bad hour costs little.\n2. **Single-node training next** — one `TaskEnvironment`, one GPU resource string, checkpoints to object storage. At this point you've validated that your image builds, your data streams, and your logs are where you expect.\n3. **Multi-node training last** — once images, data, and observability are proven. `ClusteredTaskEnvironment` is the piece with the most moving parts and the least tolerance for a half-migrated environment.\n\nNothing forces a big bang: Slurm and Flyte can run side by side indefinitely, and a Flyte task can `subprocess` out to `sbatch` during the overlap if you need a bridge.\n\n## Gotchas\n\n- **`flyte.map` returns a generator.** Wrap it in `list()` to materialize results.\n- **`memory`, not `mem`.** And GPUs use a combined `\"A100:8\"` string — type and count together, not `--gres` plus a separate accelerator argument.\n- **`shm=\"auto\"` matters for PyTorch DataLoader.** The container default `/dev/shm` is tiny; multi-worker data loading will fail cryptically without it.\n- **A clustered task cannot launch subtasks.** Orchestrate from a plain `TaskEnvironment` with `depends_on=[clustered_env]`, or you'll hit an image-cache lookup failure at runtime.\n- **`nproc_per_node` must not exceed the GPU count.** `ClusteredTaskEnvironment` validates this locally and raises `ValueError` before anything is submitted.\n- **Only rank 0 should return outputs.** Every replica runs the task body; have non-zero ranks return early or return a trivial value.\n- **Retries have no platform cap.** Total attempts equal `retries + 1`, so audit any large values ported from a `--requeue` habit.\n- **`interruptible=True` with zero retries runs on-demand.** The final attempt always falls back off spot, and a single attempt _is_ the final attempt.\n- **Clustered tasks are torchrun-focused.** Classic MPI HPC codes (`mpirun`, tightly-coupled CFD, molecular dynamics) are not the target. Keep those on Slurm.\n- **Gang scheduling and topology-aware placement are still maturing on Kubernetes.** Slurm's scheduler has decades of work behind co-scheduling N nodes on the same rack or fabric; the Kubernetes ecosystem is closing the gap but is not there.\n- **Queues order work; they don't preempt it.** `queue=\"gpu-high\"` is the nearest analogue to `--partition` plus `--qos`, but a lower-priority task already running is not interrupted. Queue names must exist in your cluster configuration, and the feature is platform-dependent — check what your deployment supports before designing around it.\n- **Frontier scale is still Slurm's.** Hundreds of GPUs per job with explicit InfiniBand topology control is where Slurm keeps winning. Migrate the pipeline layer; be deliberate about the rest.\n\n## Anti-Patterns\n\n1. **Don't recreate `/scratch` as a hardcoded bucket path.** Pass `flyte.io.File` / `flyte.io.Dir` between tasks. Two tasks agreeing on a string is exactly the coupling you're migrating away from — and it forfeits the lineage you get for free.\n2. **Don't port `#SBATCH` pragmas onto every task decorator.** Image, resources, cache, and interruptibility belong on a shared `flyte.TaskEnvironment`; override per task only where a task genuinely differs.\n3. **Don't translate `--dependency=afterok` into a sentinel-file poll.** Await the task and use its return value.\n4. **Don't keep the array-plus-collector split.** After `asyncio.gather` or `flyte.map`, reduce in plain Python inside the same task.\n5. **Don't `srun` inside a task.** The task body _is_ the rank's process. For multi-node, use `ClusteredTaskEnvironment` and let torchrun do the launching.\n6. **Don't checkpoint to a local path and expect a retry to find it.** Use `flyte.ctx().checkpoint` or write a `File` to object storage — a retry may land on a different node.\n7. **Don't do heavy compute in an orchestrating task.** A task that calls other tasks is a driver pod; CPU-bound work between awaits stalls everything downstream.\n8. **Don't lean on global state in a reusable environment.** `ReusePolicy` keeps the process alive across tasks — cache the model deliberately, and don't let mutable state leak between invocations.\n9. **Don't migrate tightly-coupled MPI simulation first (or at all).** Start with pipeline-shaped work; leave the workloads Slurm is genuinely better at on Slurm.\n\n## Related skills\n\n- **`flyte-sdk-ml`** — greenfield ML authoring in Flyte 2 (training, HPO, inference) once the migration shape is clear.\n- **`flyte-migrate`** — the entry point if you _also_ have Flyte 1 (`flytekit`) code to port.\n- **`flyte-sdk-ship`** — image specs, dependency management, and reproducible builds, i.e. everything that replaces `module load`.\n- **`flyte-sdk-app`** — serving and endpoints, for the step after training that Slurm never covered.\n"
}

SHA-256 of public snapshot: ced691fa1c0ae92b9e021a7fe0e4d6d474a086748ba4ba7ee41332344a3188c9