nemo_rl.data_plane.factory#
Single entrypoint that maps a :class:DataPlaneConfig to a client.
Module Contents#
Functions#
Whether the data plane is on. |
|
Pick the synchronous trainer based on |
|
The |
|
Set backend env vars that must be identical in every process. |
|
Construct the configured data-plane client. |
API#
- nemo_rl.data_plane.factory.data_plane_enabled(
- cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
Whether the data plane is on.
None(key absent) means off.
- nemo_rl.data_plane.factory.select_sync_trainer(
- master_config: nemo_rl.algorithms.grpo.MasterConfig,
- *,
- label: str = 'GRPO',
Pick the synchronous trainer based on
data_plane.enabled.Shared by every launcher so trainer choice cannot drift between them. Pairs with :func:
make_policy_factory: turning the data plane on means picking both the TQ-mediated trainer and aTQPolicy, and a launcher that picked only one would fail after a full model load.- Parameters:
master_config – The resolved config; only
data_planeis read.label – Algorithm name for the progress line (e.g.
"VLM GRPO").
- nemo_rl.data_plane.factory.make_policy_factory(
- cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
The
policy_factoryforsetup(), orNonefor a plainPolicy.Lives at the launcher level so the legacy trainer stays data-plane-agnostic (architectural invariant — see
tests/unit/data_plane/test_architecture_invariants.py).
- nemo_rl.data_plane.factory.maybe_configure_data_plane_env(
- cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
Set backend env vars that must be identical in every process.
Call this on the driver before
init_ray():init_raysnapshots the driver’s environment intoruntime_env["env_vars"]and hands it to every Ray worker, which are fresh processes, so the value is in place before they run any engine code.The binding constraint is not
init_raybut that this must run before anything imports the backend’s engine, which snapshots its configuration as it loads. :func:~nemo_rl.data_plane.adapters.transfer_queue_env.configure_engine_envraises rather than silently no-op’ing if that is violated.Subprocesses the driver did not spawn through Ray inherit whatever the driver’s environment held at fork, so they are covered as long as this ran first.
No-op when the data plane is disabled or the backend has no such knobs.
- Parameters:
cfg – Data-plane config, or
Nonewhen the data plane is off.
- nemo_rl.data_plane.factory.build_data_plane_client(
- cfg: nemo_rl.data_plane.interfaces.DataPlaneRuntimeConfig | None,
- *,
- bootstrap: bool = True,
- checkpointing: bool = False,
Construct the configured data-plane client.
Dispatches on the configured implementation. TransferQueue supports cross-process transfer; the local adapter keeps colocated SFT batches in one process. Raises if data_plane is disabled — the legacy trainer (
nemo_rl.algorithms.grpo.grpo_train) should be used in that case rather than a NoOp fallback here.- Parameters:
cfg – Data-plane config; must have
enabled=True.bootstrap –
Trueon the driver — bootstraps the TQ controller.Falseon worker processes — connects to the existing controller (avoids creating a second named actor).checkpointing – Prepare storage for saving or restoring data-plane state. Derived by the caller from its existing checkpoint settings and resume path. Only used at bootstrap; workers inherit the mode from TQ.
- Returns:
A configured
DataPlaneClient; wrapped in- class:
MetricsDataPlaneClientwhen observability is enabled.