nemo_rl.weight_sync.collective_weight_synchronizer#

NCCL collective weight synchronizer for non-colocated deployments.

Handles weight transfer between policy and generation workers running on separate GPU clusters using NCCL collective communication. The policy broadcasts its weights, and generation workers receive them via the established NCCL process group.

Lifecycle per sync:

  1. policy.broadcast_weights_for_collective() – send via NCCL generation.update_weights_from_collective() – receive via NCCL

  2. Verify transfer success

No offload/restore steps are needed since policy and generation run on separate GPUs with dedicated memory.

Module Contents#

Classes#

CollectiveWeightSynchronizer

Weight synchronizer using NCCL collectives for non-colocated deployments.

Functions#

_settle_before_propagating

Let every rank finish unwinding before a refit failure reaches the caller.

API#

nemo_rl.weight_sync.collective_weight_synchronizer._settle_before_propagating(futures, budget_s, what: str) → None#

Let every rank finish unwinding before a refit failure reaches the caller.

ray.get raises on the FIRST future that fails and leaves the rest running. That is fine when the caller is going to stop, and wrong when it is going to rebuild: a communicator rebuild is itself a collective, so dispatching init_collective while some ranks are still inside the old refit means they join late or not at all, and the rendezvous times out instead of coming up.

Job 6512153 measured exactly that on the reshard kill variant. Rank 0 gave up on its own deadline, the controller went straight into the recovery, and the rebuild began – line 963 of the log – two lines BEFORE rank 1’s watchdog fired at all. The surviving generation worker then spent 300s twice failing to reach a store that never came up, and the run died at 690s having done everything else right.

Bounded, and swallowing whatever the stragglers raise: they are unwinding from the same failure the caller is already holding, and replacing it with a straggler’s version would lose the diagnosis. If the budget runs out, propagate anyway – a caller stuck here would be a worse wedge than the one being recovered from.

class nemo_rl.weight_sync.collective_weight_synchronizer.CollectiveWeightSynchronizer(
policy: Any,
generation: Any,
train_cluster: Any,
inference_cluster: Any,
refit_timeout_s: Optional[float] = None,
)#

Bases: nemo_rl.weight_sync.interfaces.WeightSynchronizer

Weight synchronizer using NCCL collectives for non-colocated deployments.

Policy and generation workers run on separate GPU clusters. Weights are synchronized via NCCL broadcast over a pre-established process group.

Parameters:
  • policy – Policy object implementing ColocatablePolicyInterface.

  • generation – Generation object implementing GenerationInterface.

  • train_cluster – RayVirtualCluster for the training workers, used to obtain the master address/port and world size for collective init.

  • inference_cluster – RayVirtualCluster for the inference workers.

  • refit_timeout_s – Deadline for one refit collective. Each participating worker arms a watchdog and aborts its own communicator when it expires, which is what lets the controller rebuild over the survivors instead of blocking in NCCL forever. None disarms it entirely, so the hang protection is lost.

Initialization

sync_weights(
*,
timer: Optional[nemo_rl.utils.timer.Timer] = None,
kv_scales: Optional[dict[str, float]] = None,
) → None#
property is_stale: bool#
_desired_membership(
absent_shards: collections.abc.Sequence[int],
train_world_size: int,
) → Optional[nemo_rl.weight_sync.membership.RefitMembership]#

The membership to track, or None for a backend that owns no DP worker group.

NOT every GenerationInterface has one. vLLM and TRT-LLM do; Dynamo, Megatron and SGLang do not, and the interface declares nothing either way – so reading self._generation.worker_group assumes a vLLM shape that this synchronizer is not entitled to assume. It holds a GenerationInterface, and Dynamo reaches it through the ordinary non-colocated branch of the factory.

That assumption broke L1_Functional_Tests_Dynamo on the plain grpo.py path:

grpo.py:1747  policy_generation.weight_synchronizer.init_communicator()
AttributeError: 'DynamoGeneration' object has no attribute 'worker_group'

None means “this backend has no shards to track”, which is the same thing an unrecorded membership already meant: every reconcile falls through to “nothing to do”, exactly as it behaved before membership tracking existed. Re-admission is only meaningful where shards exist.

Deliberately NOT mirrored on the reshard synchronizer. nccl_reshard is a vLLM-only transport that REQUIRES a worker group, so a missing one there is a misconfiguration that should fail loudly rather than silently degrade.

init_communicator() → None#
_settle_budget_s() → float#

How long to let stragglers unwind: their own deadline, plus a little.

A rank that has not given up yet will do so when its watchdog fires, which is the same refit_timeout_s every rank was armed with. Without a configured deadline there is nothing bounding them, so fall back to a fixed wait rather than blocking the recovery indefinitely.

reconcile_communicator(
absent_shards: collections.abc.Sequence[int],
force: bool = False,
) → bool#

Rebuild the refit communicator over the surviving generation shards.

model_update_group spans every train and inference rank and was built once, at setup, over the full fleet. The refit is a broadcast on that group, so a missing rank blocks it forever – inside NCCL, where it produces no error and no progress while Ray still reports every actor healthy. Rebuilding without the dead ranks is what lets the run continue.

Safe for the broadcast because rank 0 is a trainer and trainers are never excluded, so the root is stable across a rebuild and each receiver still slices the same byte stream locally.

Rebuild rather than shrink/grow. The pinned NCCL exports both – checked against nccl.core.communicator.Communicator, which also exports revoke, suspend, resume and split (2.28.9 exported only shrink; uv.lock pins 2.30.7 and 2.30.4 in the dev image already has them). So this is a choice rather than a limitation: the nccl_reshard transport has to regenerate its refit plan on any membership change whatever NCCL supports, restore is dominated by the minutes an engine takes to reload, and one path shared with init_communicator is exercised by every normal run instead of only after a failure.

shutdown() → None#