In Part 1 I covered the batch size math. In Part 2 I covered DDP, FSDP, and data sharding. Those posts described a system that works. This one is about what happens when it doesn’t.

I train diffusion language models on Frontier, the exascale supercomputer at Oak Ridge National Laboratory. Frontier has over 9,400 nodes, each with four AMD MI250X GPUs (eight compute dies per node, 64 GB HBM2e per die), connected by HPE’s Slingshot-11 interconnect. It is, at the time of writing, the fastest supercomputer in the world by peak FP64 performance. It is also a shared system used by thousands of researchers, which means my 32-node training job is a tiny tenant in a building full of other people’s workloads.

That context matters. Most distributed training guides assume you own the cluster. You don’t. Nodes fail. The filesystem goes down. Your job gets preempted. The scheduler assigns you a node with a bad GPU. A system update changes a compiler flag that breaks your kernels. The question isn’t whether something will go wrong. It’s how fast you can recover when it does.

This post is about the infrastructure I built to keep training running on Frontier despite all of that.

The filesystem that kills nodes

The first lesson Frontier taught me was about its parallel filesystem: Orion, a Lustre installation. Lustre is designed for high-throughput parallel I/O. It handles large sequential reads and writes well. What it does not handle well, in my experience, is hundreds of GPU processes doing small, random-access reads simultaneously.

I mentioned this briefly in Part 2 when I discussed data sharding. But the problem goes beyond data loading. Checkpointing is worse.

Every few hundred optimizer steps, the trainer saves a checkpoint: model weights, optimizer states, scheduler state, random number generator states. For a 1B model with AdamW, a checkpoint is about 10 GB. Writing 10 GB to Lustre from one node is fine. Writing 10 GB to Lustre from 32 nodes simultaneously, while other users are also writing their checkpoints, while the filesystem is already under load from data reads: that’s when things break.

On three separate occasions, I had training jobs crash during checkpointing with node-level failures. Not application errors. The nodes themselves became unresponsive. An OLCF admin confirmed that Orion filesystem interactions were the likely cause and that a fix was being investigated. In the meantime, I needed a workaround.

NVMe: the local escape hatch

Each Frontier compute node has two NVMe SSDs: fast, node-local storage that lives and dies with the job allocation. You request it with #SBATCH -C nvme in your batch script, and it appears at /mnt/bb/$USER/. It’s ephemeral (wiped after the job ends) but blazing fast and completely isolated from the shared filesystem.

The solution: write all checkpoints to NVMe, never to Lustre. Then run a background sidecar process that copies new checkpoints from NVMe to Lustre every five minutes. Training never touches Lustre during forward/backward/optimizer steps. The sidecar runs asynchronously. If Lustre hiccups, the sidecar retries. Training doesn’t notice.

The sidecar also tracks the best checkpoint by evaluation loss, and protects it from deletion during cleanup. It keeps the three most recent checkpoints plus the best one on Lustre. Everything else is pruned to avoid filling the quota.

The staging flow per job is:

  1. Stage in: copy the latest Lustre checkpoint to NVMe on all nodes (srun cp -r)
  2. Train: --output_dir /mnt/bb/$USER/model-name (NVMe, fast)
  3. Sidecar: polls NVMe every 5 minutes, syncs new checkpoints to Lustre
  4. Job ends: NVMe is wiped, but checkpoints are safe on Lustre

Since switching to NVMe checkpointing, I haven’t had a single filesystem-related node failure. Not one.

Self-resubmitting jobs

Frontier jobs have a maximum walltime. Batch partition jobs can run up to 24 hours. My 1B model training takes roughly 30-35 hours end to end. That’s at least two jobs. The 8B model takes weeks.

The naive approach is to submit a chain of dependent jobs: sbatch --dependency=afterany:$JOB1 script.sh. But this has problems. You have to guess how many jobs you need. If a job crashes early, the rest of the chain still runs (and immediately crashes too, wasting scheduler priority). If training finishes mid-chain, the remaining jobs do nothing.

I wrote a self-resubmitting wrapper instead. Each job, near the end of its walltime, submits a new copy of itself and then exits cleanly. The new job picks up from the latest checkpoint. This continues until training completes.

The tricky part is distinguishing between four exit scenarios:

  1. Training finished: exit code 0. Don’t resubmit.
  2. Walltime approaching: the wrapper kills training 5 minutes before the Slurm walltime, saves a checkpoint, resubmits. Normal operation.
  3. Crash: a node died, RCCL timed out, something unexpected happened. Maybe resubmit, maybe not.
  4. User cancelled: someone ran scancel. Definitely don’t resubmit.

The hard case is distinguishing crash from cancel. Both send SIGTERM to the batch script. When srun launches training across 32 nodes and one node crashes, the resulting RCCL failure propagates SIGTERM to the batch script’s process group. Same signal you’d get from scancel. You can’t just trap SIGTERM and assume it was a cancel.

The solution: trap SIGTERM and set a flag, but also query scontrol show job $SLURM_JOB_ID after the process exits. If the job state says CANCELLED, a human cancelled it. If it says FAILED or NODE_FAIL, it was a crash. The two-step check is ugly but reliable.

_CANCELLED=false
trap '_CANCELLED=true' SIGTERM SIGINT

# Run training in background so we can monitor it
srun ... &
_SRUN_PID=$!

# Start a timer that kills srun before walltime
_REMAINING=$(( SLURM_JOB_END_TIME - $(date +%s) - 300 ))
( sleep $_REMAINING && kill $_SRUN_PID 2>/dev/null ) &

wait $_SRUN_PID
_EXIT_CODE=$?

if [ $_EXIT_CODE -eq 0 ]; then
    echo "Training complete. No resubmit."
    exit 0
fi

# Check if this was a user cancel
_JOB_STATE=$(scontrol show job $SLURM_JOB_ID | grep -oP 'JobState=\K\S+')
if [[ "$_CANCELLED" == "true" && "$_JOB_STATE" == *"CANCEL"* ]]; then
    echo "Job was cancelled by user. No resubmit."
    exit 1
fi

# Otherwise: walltime or crash. Resubmit.
sbatch "$0"

This is simplified. The real version generates crash reports, emails notifications on rapid crash loops, and tracks bad nodes. But the core logic is the same: run training, catch the exit, decide whether to resubmit.

The rapid crash loop problem

One failure mode I didn’t anticipate: a bad checkpoint causes a deterministic crash at load time. Training starts, loads the checkpoint, crashes. The self-resubmit kicks in, submits a new job, which loads the same bad checkpoint and crashes again. Without safeguards, this burns GPU-hours in a tight loop.

The fix: if a job crashes within 10 minutes of starting, flag it as a rapid crash and send an email notification before resubmitting. I get the email, see the pattern, and intervene manually. It’s not automatic recovery, but it limits the damage to one or two wasted allocations instead of dozens.

Node health checks and bad node tracking

About 20-30% of my 32-node jobs on Frontier encounter at least one node failure during a multi-day training run. That’s not a complaint: it’s the expected failure rate for large-scale HPC jobs on a shared system with thousands of nodes. The hardware is being pushed hard by many users simultaneously. Failures are normal. The question is how you respond.

Pre-flight health checks

Before training starts, every allocated node runs a quick health check:

  1. GPU count: does rocm-smi report all 8 compute dies? I’ve been assigned nodes where one MI250X was in a degraded state, exposing only 6 dies.
  2. ECC errors: are there uncorrectable memory errors on any GPU? If yes, this node will corrupt tensors.
  3. Network: is the Slingshot interface up? A node with a down NIC will cause RCCL timeouts on every collective.
  4. NVMe: can we write to the local SSD? If the NVMe is full from a previous user’s data, checkpointing will fail.

The whole check runs in 30 seconds via srun across all nodes. If any node fails, the job exits immediately. No wasted walltime waiting for a RCCL timeout 20 minutes into training.

Bad node tracking

When a node fails during training, I log its hostname to a persistent file. The next job submission reads this file and passes --exclude=node001,node042,... to sbatch. The bad node list accumulates over time. Entries older than 14 days are automatically cleaned: the node has probably been serviced by then.

# .bad_nodes.txt
node1042  2026-03-18  RCCL_timeout
node0887  2026-03-19  health_check_gpu_count
node1201  2026-03-20  NODE_FAIL

There’s a known false positive issue: when you scancel a job, the RCCL shutdown messages in stderr contain hostnames of all 32 nodes, not just the one that failed. The crash parser records all of them as “bad.” After an intentional cancel, I clear the file manually. It’s annoying but not worth over-engineering a fix for.

Crash diagnostics

When a job crashes (not a cancel, not a walltime exit), the wrapper generates a crash report:

=== Crash Report: Job 4229908 ===
Job name:     mamba-1b-32n
Exit code:    1
Runtime:      2h 14m 38s
Failed nodes: node1042 (parsed from RCCL error)
Error pattern: RCCL timeout (heartbeat)
Last stderr:
  [rank 128] RCCL WARN Heartbeat timeout at ...
  [rank 128] RCCL ERROR Task 0 has timed out ...

The report is emailed automatically. I can see at a glance whether it was a transient node failure (resubmit and exclude the node) or something systemic (a bad checkpoint, a code bug, a filesystem issue).

The IPv6 debugging session I’ll never forget

Let me tell you about the two hours I lost to errno:97.

I had a 32-node job that worked perfectly on 8 nodes but hung at startup on 32. The logs showed:

[rccl] Bootstrap: Using hsn0:fd00:...  (IPv6)
[rccl] WARN Connect to fd00:... failed: errno:97
        (Address family not supported by protocol)

RCCL was resolving the master hostname to an IPv6 address. On 8 nodes, the master happened to resolve to IPv4. On 32 nodes, a different node was assigned as master, and it resolved to IPv6 first. The Slingshot CXI provider doesn’t support IPv6 for RCCL connections.

The fix:

export NCCL_SOCKET_FAMILY=AF_INET
export GLOO_SOCKET_IFNAME=hsn
export TORCH_NCCL_SOCKET_IFNAME=hsn

Force IPv4 everywhere. And as a hard safety net, validate the master address before launching:

MASTER_ADDR=$(getent ahostsv4 "$MASTER_HOSTNAME" | head -n 1 | awk '{print $1}')

if ! echo "$MASTER_ADDR" | grep -qE '^[0-9]+\.[0-9]+\.[0-9]+\.[0-9]+$'; then
    echo "FATAL: MASTER_ADDR is not IPv4: ${MASTER_ADDR}" >&2
    exit 1
fi

Three lines of validation that save hours. The frustrating thing is that this only manifests at scale. Small jobs work fine because you get lucky with DNS resolution order. You only discover the bug when you scale up and a different node becomes master. Classic Heisenbug.

The Triton compiler crash

This one was my favorite, in the way that a puzzle is your favorite after you’ve solved it.

Mamba’s backward pass uses custom Triton kernels for the selective scan operation. On Frontier’s MI250X (AMD gfx90a architecture), these kernels compiled and ran correctly for months. Then I updated to Triton 3.6.0 and everything broke:

RuntimeError: PassManager::run failed
  in _chunk_state_bwd_db_kernel

The forward pass was fine. Only the backward pass crashed. And only on gfx90a: the same code ran without issues on MI300X (gfx942).

After a lot of digging through Triton’s MLIR pipeline, I found the culprit: a new optimization pass called ConvertToBufferOps. It rewrites memory operations to use hardware buffer descriptors instead of raw global loads. On gfx942 (CDNA3), this works correctly. On gfx90a (CDNA2), the pass hits an out-of-bounds index when analyzing the complex pointer arithmetic in Mamba’s backward kernels.

The fix is two environment variables:

export AMDGCN_USE_BUFFER_OPS=0
export AMDGCN_USE_BUFFER_ATOMICS=0

This tells Triton to skip the buffer descriptor optimization and fall back to plain global_load/global_store. Performance impact: negligible. Mamba’s backward kernels are compute-bound, not memory-bound. The buffer descriptor optimization would save a few percent on memory latency, but the kernels spend most of their time in FMA operations anyway.

I now set these variables in every Frontier training script. Non-negotiable. The pass might get fixed in a future Triton release, but I’m not going to find out by having my training crash at hour 20 of a 24-hour job.

Offline everything: the compute node island

Frontier’s compute nodes have no internet access. None. This sounds like a minor inconvenience until you realize how many things quietly phone home.

HuggingFace’s AutoConfig.from_pretrained() tries to download config files. Wandb tries to sync runs to the cloud. Pip sometimes checks for package updates during import. Any of these will hang for 30 seconds, time out, and either crash or silently degrade.

The environment variables that make everything work offline:

export HF_HUB_OFFLINE=1
export TRANSFORMERS_OFFLINE=1
export HF_HOME=/lustre/.../hf_cache
export WANDB_MODE=offline

Wandb in offline mode writes run data to a local directory. I sync it to the cloud manually from a login node (which does have internet) after the job finishes. This introduces its own complication: self-resubmitting jobs create a new wandb run directory every time they restart. If you don’t set a fixed WANDB_RUN_ID, you end up with 15 fragmented runs instead of one continuous training curve.

export WANDB_RUN_ID="mamba-1b-32n-v2"    # deterministic, not random
export WANDB_RESUME=allow                  # append to existing run

These two lines took me a full day to figure out. I had beautiful loss curves that abruptly ended and restarted every 6 hours (one walltime cycle). The data was all there, split across separate runs. Stitching them together manually is possible but miserable. Setting a fixed run ID solves it permanently.

MFU: the number everyone asks about

Model FLOP Utilization (MFU) is the fraction of the GPU’s theoretical peak compute that your training actually uses. It’s the single most important efficiency metric for large-scale training. If your MFU is 5%, you’re wasting 95% of the hardware. If it’s 40%, you’re doing well.

The formula:

\[\text{MFU} = \frac{\text{FLOPs per step}}{\text{step time} \times \text{peak TFLOPS} \times 10^{12}}\]

For a 1B Mamba model on MI250X:

\[\text{FLOPs/token} \approx 6N = 6 \times 10^9\] \[\text{FLOPs/step} = 6N \times \text{tokens\_per\_step} = 6 \times 10^9 \times 16{,}384 \approx 98.3 \times 10^{12}\]

Note: that’s per GPU. Each GPU processes 4096 × 1024 / 256 = 16,384 tokens per step. MFU is a per-device metric.

\[\text{MFU} = \frac{98.3 \times 10^{12}}{3.5 \times 95.74 \times 10^{12}} \approx 29.3\%\]

Each MI250X compute die has a peak of 95.74 BF16 TFLOPS. At ~29% MFU, we’re getting about 28 TFLOPS of useful compute per die. The rest is overhead: memory bandwidth bottlenecks, communication latency, Python overhead, kernel launch gaps, and the sequential nature of Mamba’s selective scan (which doesn’t parallelize as well as attention’s matrix multiplications).

For reference, the well-optimized attention-based LLM training runs (like the ones reported in the Llama papers) achieve 35-45% MFU on NVIDIA hardware. Getting 30% on AMD hardware with a Mamba architecture (which has inherently sequential SSM scans) is reasonable. There’s room to improve, but it’s not embarrassing.

The $6N$ approximation for FLOPs per token is a rough estimate. The actual count, computed by walking the model’s modules and summing per-layer FLOPs, is 52.09 billion FLOPs per token for the 1B BiMamba model (99.5% from linear layers, 0.5% from SSM operations). The SSM FLOPs are tiny because they scale with state dimension, not sequence length. Attention FLOPs, by contrast, scale quadratically with sequence length: for the equivalent attention model, attention score computation accounts for 3.5% of total FLOPs.

What I would do differently

Looking back at nine months of training on Frontier, a few things stand out.

Start with NVMe checkpointing from day one. I lost two weeks of cumulative training time to Lustre-related crashes before switching. The sidecar pattern is simple to implement and eliminates an entire class of failures.

Validate the environment before training starts, not after it crashes. The health checks, IPv4 validation, and offline mode flags should be in every script from the beginning. I added most of them reactively, after each bug cost me hours of debugging.

Fix your wandb run IDs immediately. Fragmented training curves are not just an inconvenience. They make it genuinely hard to compare runs, spot regressions, and decide when to stop training. The cost of getting this wrong is invisible until you’re staring at 15 disconnected line segments trying to figure out if your loss actually converged.

Keep a bad node list. Frontier has thousands of nodes and some of them are flaky. Excluding known-bad nodes from your submissions is the cheapest reliability win available. It takes 20 lines of bash.

None of this is novel computer science. It’s plumbing. But plumbing is what determines whether your 32-node job produces a trained model or a collection of crash logs.

Acknowledgments

This research used resources of the Oak Ridge Leadership Computing Facility (OLCF) at Oak Ridge National Laboratory, which is a DOE Office of Science User Facility supported under Contract DE-AC05-00OR22725. Frontier is operated by the OLCF and was the first exascale computing system, achieving 1.1 exaflops on the HPL benchmark.

I’m grateful to the OLCF support team for their responsiveness on the filesystem issues and for maintaining a system that, despite the war stories in this post, is remarkably stable for its scale. The node failure rate I reported is normal for jobs spanning hundreds of nodes on any large HPC system: it’s a property of the scale, not a deficiency of the machine.

References

Frontier and OLCF:

  • Frontier User Guide. The official documentation for running jobs on Frontier: Slurm configuration, module system, filesystem layout, and NVMe burst buffers.
  • OLCF Frontier System Overview. Architecture specs: 9,408 nodes, AMD EPYC 7A53 CPUs, MI250X GPUs, Slingshot-11 interconnect.

Distributed training infrastructure:

Hardware and interconnects:

  • AMD, MI250X Datasheet. The GPU specs: 110 CUs per GCD, 128 GB HBM2e per package (64 GB per die), 95.7 TFLOPS BF16 per GCD.
  • Hewlett Packard Enterprise, Slingshot Interconnect. The network: 200 Gb/s per NIC, 4 NICs per node, dragonfly topology.

Mamba and SSMs:

MFU and training efficiency:

Scaling laws: