latest (dev)
Copy
Latest development documentation · Updated 2026-10-08
tensorplay.distributed.elastic
tensorplay.distributed.elastic runs a training job across a set of nodes
that can change size while the job is live. The mental model has three moving
parts:
A rendezvous decides who is in the world. Each node asks a shared rendezvous store “how many nodes and which roles have joined?” and the answer changes over time: nodes can join and leave between training steps, so the world size and per-rank world mapping are renegotiated rather than fixed at launch.
An agent runs on every node. It reads the rendezvous result, launches the local workers (processes for a subprocess entrypoint, workers for a callable), monitors their health, restarts them on failure up to
max_restarts, and re-rendezvous when the peer set changes.The workers are your actual training processes. The elastic agent gives each one the standard environment variables (
RANK,WORLD_SIZE,MASTER_ADDR,MASTER_PORT,LOCAL_RANK,LOCAL_WORLD_SIZE) so the training script initializes the distributed package as it would under any launcher.
The entry point is main(), the
python -m tensorplay.distributed.run command’s implementation, which parses
the CLI into a LaunchConfig and
hands the job to the elastic agent.
# train.py -- runs on every worker
import tensorplay
import tensorplay.distributed as dist
dist.init_process_group("nccl", rank=dist.get_rank(), world_size=dist.get_world_size())
# ... training loop ...
python -m tensorplay.distributed.run \
--nproc-per-node=2 --nnodes=1 \
--rdzv-backend=static --rdzv-endpoint=localhost:29500 \
train.py
When a worker process later calls dist.monitored_barrier/dist.barrier,
the elastic agent can observe a stuck or dead worker and either restart it or
fail the job, depending on the remaining restart budget.
The launcher
Elastic launch settings for one job. |
|
Callable wrapper around |
|
Start the local agent for |
LaunchConfig carries the job
definition: the min/max node count (min_nodes/max_nodes), workers per node
(nproc_per_node), the rendezvous backend and endpoint, the restart and
monitor settings, and the start method (spawn by default). When min_nodes
and max_nodes differ the job is elastic and ranks can be added or removed
while it runs. elastic_launch() wraps
launch_agent and is the programmatic form of the CLI: build a
LaunchConfig, pass your entrypoint callable or script path, then call the
resulting function with any trailing CLI arguments.
The agent
Blueprint of the worker group this agent manages. |
|
One logical worker slot with its rank assignments. |
|
Mutable group state driven by the agent. |
|
State of the worker group in the agent run loop. |
|
Terminal outcome of the agent run for one role. |
|
Agent interface for one worker-group role. |
|
|
Reusable agent run loop over one worker group. |
|
Agent managing workers on the local node. |
WorkerSpec is the
blueprint of the local worker group: the role name, local_world_size, the
entrypoint (fn callable or entrypoint command plus args), and the
settings that bound restarts (max_restarts), the rendezvous handler, and how
workers are launched (start_method, redirects, tee, log_dir, and the
environment variables derived from the spec). Every node runs the same spec,
so the resulting world arithmetic is consistent across nodes.
ElasticAgent is the
abstract agent interface over one worker-group role. The
SimpleElasticAgent run
loop is the reusable implementation: it alternates between rendezvous (getting
the current membership) and worker monitoring, restarting workers that exit
with a non-terminal state and re-rendezvousing when other nodes are waiting.
LocalElasticAgent
manages the workers on one node, launching them through start_processes
(with per-rank environments) and letting the base-class loop handle failures
and restarts. WorkerState
is the lifecycle of a single worker slot (INIT, HEALTHY, UNHEALTHY,
SUCCEEDED, FAILED, STOPPED); RunResult
records the terminal outcome per role.
Rendezvous
Algorithmic interface of one rendezvous backend. |
|
|
Parameters describing one rendezvous request. |
Outcome of a successful rendezvous for one agent. |
|
|
Connection information for the bootstrap store handed to workers. |
|
Configuration shared by the rendezvous state machine. |
|
Coordinate membership, ranks, heartbeats, and round transitions. |
|
One-shot rendezvous over a static TCPStore endpoint. |
|
Rendezvous state backend backed by a TP key/value Store. |
|
Store rendezvous state behind compare-and-set operations. |
Create a handler from the process-wide registry. |
RendezvousHandler is the
algorithmic interface behind one rendezvous backend. The main method is
next_rendezvous(),
which blocks until the world reaches a consistent membership and returns a
RendezvousInfo describing
who joined and their rank-to-world mapping. The handler also reports how many
nodes are waiting (num_nodes_waiting) which the agent uses to trigger
scale-up, and supports closing and shutting down the rendezvous.
RendezvousParameters
describes one rendezvous request: the backend name, the store endpoint,
the run_id, the min_nodes/max_nodes bounds, and backend-specific options
in config. StaticTCPRendezvous
implements the simplest backend over a static TCP store (fixed membership),
DynamicRendezvousHandler
is the dynamic variant used by the etcd/tcp backends that allow membership to
change, while P10dRendezvousBackend
and TpRendezvousBackend
run over the key/value stores the process group already provides
(TCPStore, FileStore), so a rendezvous needs no third-party service.
create_handler() builds a
handler from a backend name and parameters.
Base class for all rendezvous failures. |
|
|
The rendezvous has been closed and accepts no more participants. |
|
A rendezvous phase exceeded its time budget. |
|
The rendezvous backend store could not be reached. |
|
The rendezvous state is corrupt or the backend rejected an update. |
|
Raised to unwind an agent when it is the last node leaving. |
|
The rendezvous join window elapsed without reaching min nodes. |
The rendezvous error types let a caller distinguish why membership could not be formed: the rendezvous was closed, it timed out, a connection was lost, the world reached an inconsistent state, it was asked to exit gracefully, or the maximum node count was exhausted.
Multiprocessing
|
|
Base class owning a homogeneous group of worker processes. |
|
|
Outcome of monitoring a worker group to completion. |
|
Structured failure of one worker process. |
|
Raised by the launcher when one or more workers failed. |
|
|
Per-stream redirection modes for a worker group. |
|
Which standard streams a worker's output should go to. |
|
Expand a per-rank redirection spec into one entry per local rank. |
start_processes()
launches len(envs) workers and returns the managing
PContext. The
entrypoint is either a command string (subprocess workers, where args is a
shared argument list) or a picklable callable (multiprocessing workers, where
args holds one tuple per rank). Output redirection and logging are configured
through Redirects and
Std, so a worker’s
stdout/stderr can be teed to the log directory or suppressed. When a child
fails, the exception that propagates to the manager is a
ChildFailedError (for
multiprocessing workers) or a
ProcessFailure (for
subprocesses) carrying which worker failed and its output.
Data utilities
Wrap a finite iterable into an endless one. |
|
|
Sampler for elastic jobs where |
|
Return the value of |
Create a socket bound to an ephemeral local port and return it. |
|
Substitution variables usable in worker argument templates. |
ElasticDistributedSampler
is the data-parallel sampler you use when the world size can change: it
shards the dataset across the current ranks and, crucially, produces a
world-sized sample stream that stays stable as nodes join or leave, so the
resumed epochs do not reorder the data. CyclingIterator
wraps an iterator and repopulates it from a generator function whenever it is
exhausted, so a training loop can run for an arbitrary number of steps even
when a node’s local dataset is finite. get_env_variable_or_raise()
reads a required environment variable and raises a helpful error when missing,
macros is the collection of
elastic-injected macro functions (such as the world-size and local-rank
helpers) that the agent makes available to workers.
Events, metrics, and timers
A single agent or worker lifecycle event. |
|
Lifecycle state carried by events. |
|
A rendezvous progress event, rendered as JSON by the logging handler. |
|
Dispatch |
|
Route |
|
Metric plumbing configuration; maps groups to handler names. |
|
Decorator measuring wall time of |
|
Client side of the timer contract. |
|
Watches outstanding deadlines and reacts when they expire. |
|
Set the process-wide default |
|
Context manager asserting the block finishes within |
The instrumentation is pluggable so the same agent loop can emit events to a
logging handler, publish metrics to a stream, or enforce deadlines.
Event describes one agent or
worker lifecycle event; NodeState
is the set of node lifecycle states, and record()
dispatches an event to the handler configured for a destination.
MetricsConfig /
configure() select where metrics
go, and prof() times a block and
publishes the duration. TimerClient
sends deadline requests to a TimerServer,
letting a worker register a “this step must finish before X” deadline;
expires() is the context manager
that raises if the deadline is not met.
Where to go next
the distributed package — process groups, collectives, and initialization, which the launched workers call in their training loop.
device mesh — multi-dimensional group layout for tensor and FSDP parallelism once the elastic world is formed.
FSDP — the sharding strategy for the model the elastic job trains.
Help improve this page
Found an error, an unclear step, or a missing example?
tensorplay.distributed.device_mesh
A device mesh is the execution context for distributed tensors . It is an n-dimensional array whose entries are global ranks: the value at coordinates (i, j, ...) is the rank of the process holding that position of the m
tensorplay.distributed.fsdp
Fully sharded data parallelism (FSDP) shrinks a model’s peak GPU memory by sharding its parameters across the ranks of a process group. The classic wrapper flattens a group of parameters into a single buffer and splits i

