Optimizing Effective Training Time for Large-Scale Recommendation Systems
Mingming Ding, Ruilin Chen, Yuzhen Huang, Hang Qi, Menglu Yu, San Tan, Damian Reeves, Boris Sarana, Kevin Tang, Satendra Gera, Gagan Jain, Sahil Shah, Vishwa Karia, Fuzail Khan, Yashasvi Makin, Edward Z. Yang, Oguz Ulgen, Jia Chen Ren, Laith Sakka, Mayank Garg, Meet Vadakkanchery, Aici Lin, Wei Sun, Mengjiao Zhou, Shuai Yang, Junqing Zhou, Max Leung, Apoorv Purwar, Musharaf Sultan, John Bocharov, Zhenyu Tang, Vivek Trehan
cs.IR
2026-10-02
Meta turned training-time loss into an owner-mapped ETT% metric, then fixed init, compilation, checkpoints and publishing; six models gained 15.5 points, the fleet topped 90%.
Meta's largest ads recommendation models train on tens of billions of examples a day across thousands of GPUs, and training is continuous: a new data window, a hyperparameter change, an infrastructure failure, or a preemption each forces a restart. A job holds its GPUs from allocation to release, and before this work only 50-60% of end-to-end wall time advanced training on new data. The rest went to initialization, compilation, checkpointing, publishing, and recovery.
Existing metrics cannot see this loss. MFU covers the training loop only; a job can sustain 95% MFU on every step while spending just 60% of its wall time on new data. Google's Goodput covers the full lifecycle but does not decompose: a compiler regression, a scheduler change, and a storage slowdown all move the top-level number the same way, and no team can act on it. Recommendation training amplifies the gap structurally. Jobs are short: a 51-minute cold start is 34% of a 151.5-minute job and 0.2% of a three-week LLM run. Architectures mix huge embedding tables with dynamic shapes, so PT2 compilation can exceed an hour from a cold cache. Streaming data adds a 3-5 minute warm-up before the first step.
The core contribution is an attributable metric; the individual optimizations are established systems methods, and the paper's claim is the process for choosing them, measuring end-to-end return, and holding the gains.
ETT% is the fraction of wall time spent consuming new data, written as 1 - (TTS + NoF × TTR) / wall time. TTS is per-start cold-start idle, TTR is per-failure recovery idle, NoF is the failure count. A second level splits this into six components (scheduling, trainer init, PT2 compilation, effective training, wasted training, shutdown), each mapped to one infrastructure system and one owning team, so a dashboard regression names the team that can fix it.
| Setting | Metric | Result |
| 6 paired benchmarks (8-256 GPUs, 2 induced failures each) | ETT% | 59-85% baseline to 80-93%, mean +15.5 points |
| Fleet, 500+ models, 6 months | ETT% | 80% to above 90% |
| Cold-cache PT2 compile, 32 H100s | Compile time | 2,050s to 254s (-87.6%) |
| Models with standalone publishing | Shutdown time | 25.6-54.6 min to 3.3 min or less |
| Largest workload | ETT% | 85% |
Recovery accounts for 58% of the time saved. A recovery repeats the startup phases, so TTR fell by 1.8-8.6 times the TTS reduction on every model. Direct attribution credits the trainer-start work with 8% of savings; the restart-cascade bound puts it at 24.8%, a factor of three. A single top-level scalar systematically undervalues startup work.
The transferable part is the operating model: a two-level decomposition where every component has an owner and a service-level alert, run as a dashboard that held the gains for six months while models and infrastructure churned. The individual techniques have prior art, as the authors say themselves, and several fixes shipped in open-source TorchRec and PyTorch (DCP AsyncStager, Mega-Cache), so they are usable outside Meta. For three-week LLM pretraining runs this barely matters, with cold start at 0.2% of wall time. For recommendation, continuous learning, and any short-job fleet on shared clusters, the arithmetic is concrete: at 1,000 GPUs, cold start plus restarts cost about 3,400 GPU-hours a day before the fixes.
The authors list four: one run per configuration, so single measurements without confidence intervals (repeating baseline pairs at up to 2,000 GPUs was unaffordable); standalone publishing was not active on every benchmark model; shared clusters add scheduling variance that was not isolated; and no leave-one-out ablation, so the trainer-start contribution is a bound, not an isolated measurement. The fleet record has its own caveats: models migrated platforms mid-window and the fleet population changed, so the trend reflects the fleet as operated.
Two things stand out beyond the self-report. The text says seven models were evaluated while Figure 6 lists six (A-F); the paired controlled evidence tops out at 256 GPUs, so the 2,000-GPU scale and the 85% figure for the largest workload rest on the longitudinal dashboard. And each paired job is killed exactly twice, which fixes NoF for attribution but does not reflect the natural failure mix.