Research
Nereus: Adaptive Parallelism for LLM Post-Training
Overview Research area: Distributed systems for machine learning — specifically online adaptation of parallel execution plans during reinforcement learning (RL) post-training of large language models

- arXiv
- 2609.34645
- Published
- 2026-09-28
- Authors
- Songlin Jiang, Tuo Shi, Sitong Zhang, Zeke Wang, Mario Di Francesco, Bo Zhao
AI summary
Overview
Research area: Distributed systems for machine learning — specifically online adaptation of parallel execution plans during reinforcement learning (RL) post-training of large language models (LLMs).
Technical level: Advanced. The paper assumes familiarity with data, tensor, and pipeline parallelism (DP/TP/PP), collective communication, GPU memory budgeting, and RL post-training workflows such as PPO.
Scope: The paper presents Nereus, a cost-aware runtime that replans and transitions a running multi-model RL post-training job across changing GPU availability, sequence lengths, and bottlenecks, while reusing GPU-resident state rather than restarting.
What This Paper Is About
RL post-training for LLMs coordinates several models (actor, critic, reward, reference) across generation, inference, and training stages on shared GPUs. During a run, the GPU pool can shrink or grow, generated sequences get much longer, memory pressure rises, and the bottleneck can move from one stage to another — so the execution plan chosen at startup can become slow or impossible. Nereus's goal is to detect this "drift" and move the job to a better plan online, reusing the existing distributed state and paying only a small, bounded transition cost.
Key Contributions
- A low-overhead, cost-aware adaptation policy. A controller combines event-driven triggers with an online-calibrated cost model, selects a memory-feasible global plan, and admits a transition only if the predicted step-latency savings repay the estimated transition cost (the payback rule in Eq. 2). It also replans from intact units after node failures.
- A dependency-aligned state abstraction, the Elastic Model Unit (EMU). An EMU is one model-stage replica that encapsulates TP/PP and its shards internally, while exposing DP replicas externally for scaling. The boundary matches the job's coupling structure: TP/PP are tightly synchronized and stay inside the unit, while DP replicas are loosely coupled.
- Four core state-transformation primitives. Split, Merge, Extend, and Destroy cover parallelism resharding and resource scaling across the plan space while preserving training semantics and logical state.
- Safe concurrent transition orchestration. The engine compiles the difference between the current and target plans into a global transition DAG, adding cross-model-stage resource-dependency edges when transient GPU overlap blocks an acquisition, and prioritizing resource-releasing operations to avoid deadlock.
Main Findings
- Throughput gains over existing frameworks: Nereus improves end-to-end 8B PPO throughput by 2.14–7.27× (median 3.99×) over OpenRLHF and by 1.10–1.47× over Verl across diverse clusters.
- Near-optimal plan selection: Its selected plans stay within 5% of the empirical optimum in all 18 measured settings.
- Online TP/PP adaptation pays off: In a trace built from real data, online TP/PP adaptation reduces average step latency by 27.7% relative to the initial fixed TP/PP layout with DP scaling.
- Transitions are cheap relative to the run: In a 1,000-step run reaching 1,024 GPUs, six transitions consume 0.079% of total run time.
- Better than admission without payback control: Across three held-out traces, Nereus averages 858.7 s/step compared to 928.3 s for DynaRL-style admission.
- EMU abstraction scales faster than prior state boundaries: Scaling with EMUs is 3.8–16.2× faster than with Oobleck and Tenplex.
- Coordinated transitions are more reliable under overlap: Coordinated transitions succeed in all overlap trials, compared to only 34–62% for DynaRL's per-component migrations.
- The model-stage boundary is the cheapest of the three compared: For an 8B actor/critic workload scaled from 16 to 32 GPUs on Cluster #3, the measured cost is 836.74 s for checkpoint/restart, 66.43 s for shard-level state management, and 6.52 s for the model-stage (EMU) boundary.
- Drift is large and measured: Generated sequence length grows from 500 to 8,000 tokens within 1,000 steps for Llama-3.1-8B, a 16× increase; at 16 GPUs, the best (DP, TP, PP) training plan changes with this growth, where (4,2,2) is fastest at 2K tokens but runs out of memory at 4K, while (4,4,1) stays feasible at 4K but adds communication overhead at 2K.
- Resource supply is volatile: An analysis of AWS traces found more than ten availability changes for a 16-GPU job over ten hours, and in plan enumeration for Llama-3.1-8B the best TP/PP layout differs between 32 and 256 GPUs.
- Offline predictors are unreliable: For Llama-3.1-8B, an offline predictor selects an infeasible plan at 8 GPUs, and at 64–256 GPUs its selected plans are up to 1.56× slower than the best measured plans — motivating online calibration.
Methodology in Plain English
Nereus splits the problem into a control loop and an execution engine.
The monitor watches sequence length, available GPUs, peak memory, and achieved compute and communication efficiency using Ray and NVML (or ROCm SMI) and framework instrumentation. It replans when the GPU pool changes, when a memory violation is predicted, or when a drifting signal deviates from its last reference value by more than a relative threshold (30% for per-step sequence length in all experiments). Reference values reset after every replan so steady growth keeps triggering.
The planner enumerates candidate DP/TP/PP configurations per model-stage, estimates each one's compute, collective, pipeline-bubble, and memory costs, and uses dynamic programming to pick the lowest predicted end-to-end step latency under a per-GPU memory cap and stage-wise GPU budgets. Instead of relying on peak hardware rates, it learns one efficiency correction per operation class (such as a GEMM, attention, or a specific all-reduce on a link type) from measured execution times, so the cost model is calibrated against the running job rather than a full offline profile. Replanning runs in a single CPU process on the head node and uses no GPUs.
The admission rule applies two criteria. Feasibility checks that GPU memory can hold retained state, workspaces, intermediate units, co-resident model state, and collective buffers (with 10% headroom). Profitability compares predicted per-step savings against the estimated transition cost, admitting only when the payback period fits within a confidence fraction (γ = 0.5 by default) of the estimated time until the triggering signal next crosses its threshold. Resource revocation follows an urgent path that waives profitability and can fall back to checkpoint/restart.
The transition engine represents current and target plans as collections of EMUs and reduces each model-stage's change to a local DAG in the order Split, Destroy, Extend, Merge. It then adds cross-model-stage edges only where transient GPU overlap blocks an acquisition, using release-before-acquire prioritization and generation-to-inference-to-training tie-breaking. Ready primitives run concurrently across GPU groups and streams, with Split using position-wise AllGather, Merge assembling from replicated state, Extend using rank-aligned Broadcast, and Destroy releasing resources.
Nereus is implemented in 43k lines of Python, C++, and CUDA/HIP, integrating with vLLM, DeepSpeed, and Megatron-LM over standard NCCL/RCCL collectives.
Why This Matters
Impact on research. The paper reframes elastic execution for RL post-training as a question of choosing the right state boundary, not just the right parallelism degrees. Its comparison of checkpoint-level, shard-level, and model-stage boundaries — 836.74 s versus 66.43 s versus 6.52 s on the same 16-to-32-GPU scaling workload — gives the field a concrete measurement of how much abstraction choice costs. The payback-based admission rule extends prior single-job ideas (Pollux, Sia) to coupled multi-model workloads.
Real-world applications:
- LLM alignment pipelines using PPO, ReMax, or GRPO, where long-context generation changes memory and throughput demands mid-run.
- Cloud and spot-instance training, where GPUs are reclaimed or reallocated and availability changes more than ten times over ten hours for a modest job.
- Multi-tenant GPU clusters, where network contention and interference shift which parallel layout is fastest even when allocation and workload stay fixed.
- Long-running post-training jobs, where a single run can consume 100,000 GPU-hours and restarting to reshard state is prohibitively expensive.
Industry relevance. The 2.14–7.27× PPO throughput improvement over OpenRLHF and 1.10–1.47× over Verl come from the same hardware, so the gain is a scheduling and state-management gain rather than a hardware gain. Systems that cannot change plans after startup must either over-provision for worst-case sequence length or risk out-of-memory failures — Nereus removes that trade-off.
Future Directions
- Generalizing the EMU abstraction beyond the evaluated RL setup. The paper covers 8B actor/critic PPO jobs; whether the model-stage boundary remains the right unit for mixture-of-experts models, multimodal training, or different RL algorithms is an open question.
- Tightening the drift-horizon estimate. Admission depends on H_x, an exponential moving average of intervals between threshold crossings, which starts at one step before any interval is observed. More robust forecasting or sensitivity analysis beyond the reported γ evaluation would strengthen admission decisions.
- Coupling adaptation more deeply with cluster schedulers and fault-tolerance systems. Nereus replans from intact units after GPU loss and relies on external fault-tolerance systems for state whose last copy is lost; tighter integration is a natural next step.
- Extending the transition scheduling comparison. The paper compares its DAG construction against an off-the-shelf SCIP solver on the scheduling problem; broader comparisons against additional elastic-training and RL scheduling baselines would clarify where the gains come from.
Target Audience
Systems and machine learning infrastructure researchers and engineers who build or operate distributed LLM post-training stacks — particularly those working on parallelism planning, elastic training, checkpointing and state migration, or cluster scheduling. It is also relevant to practitioners running RLHF-style pipelines on shared or preemptible GPU clusters, and to graduate students studying the systems side of large-model training. Readers need comfort with DP/TP/PP terminology and collective communication to follow the design, though the high-level problem framing is accessible.
Authors’ abstract
Reinforcement learning (RL) post-training for large language models (LLMs) coordinates multiple models across generation, inference, and training on GPU clusters. Several factors may change during a run, including resource availability, sequence length, memory pressure, and stage bottlenecks. As a consequence, an execution plan that was initially suitable can then become slow or even infeasible over time. However, adapting a job whose models share GPUs entails significant challenges: deciding whether a new plan is worth the transition cost, reusing the job's distributed state, and coordinating GPU transfers across models and stages. Nereus targets these challenges as a cost-aware runtime that adapts RL post-training jobs into efficient execution plans. Its low-overhead controller selects a memory-feasible global plan and admits the transition using a cost model calibrated against the running job. To estimate and execute a transition, Nereus represents the distributed state of each replica of a model-stage (one model in one stage) as an Elastic Model Unit. It then employs a global transition graph to order the transformations and GPU transfers of these units. In a trace built from real data, online TP/PP adaptation reduces average step latency by 27.7% relative to the initial fixed TP/PP layout with DP scaling. In a 1,000-step run reaching 1,024 GPUs, six transitions consume 0.079% of total run time. Nereus improves end-to-end 8B PPO throughput by 2.14--7.27$\times$ over OpenRLHF and by 1.10--1.47$\times$ over Verl across diverse clusters.