TensorPlay
Reference guides
latest (dev)
Copy
View Markdown

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

tensorplay.distributed.launcher.LaunchConfig

Elastic launch settings for one job.

tensorplay.distributed.launcher.elastic_launch

Callable wrapper around launch_agent().

tensorplay.distributed.launcher.launch_agent

Start the local agent for entrypoint and run it to completion.

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

tensorplay.distributed.elastic.agent.server.WorkerSpec

Blueprint of the worker group this agent manages.

tensorplay.distributed.elastic.agent.server.Worker

One logical worker slot with its rank assignments.

tensorplay.distributed.elastic.agent.server.WorkerGroup

Mutable group state driven by the agent.

tensorplay.distributed.elastic.agent.server.WorkerState

State of the worker group in the agent run loop.

tensorplay.distributed.elastic.agent.server.RunResult

Terminal outcome of the agent run for one role.

tensorplay.distributed.elastic.agent.server.ElasticAgent

Agent interface for one worker-group role.

tensorplay.distributed.elastic.agent.server.SimpleElasticAgent

Reusable agent run loop over one worker group.

tensorplay.distributed.elastic.agent.server.LocalElasticAgent

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

tensorplay.distributed.elastic.rendezvous.RendezvousHandler

Algorithmic interface of one rendezvous backend.

tensorplay.distributed.elastic.rendezvous.RendezvousParameters

Parameters describing one rendezvous request.

tensorplay.distributed.elastic.rendezvous.RendezvousInfo

Outcome of a successful rendezvous for one agent.

tensorplay.distributed.elastic.rendezvous.RendezvousStoreInfo

Connection information for the bootstrap store handed to workers.

tensorplay.distributed.elastic.rendezvous.RendezvousSettings

Configuration shared by the rendezvous state machine.

tensorplay.distributed.elastic.rendezvous.DynamicRendezvousHandler

Coordinate membership, ranks, heartbeats, and round transitions.

tensorplay.distributed.elastic.rendezvous.StaticTCPRendezvous

One-shot rendezvous over a static TCPStore endpoint.

tensorplay.distributed.elastic.rendezvous.P10dRendezvousBackend

Rendezvous state backend backed by a TP key/value Store.

tensorplay.distributed.elastic.rendezvous.TpRendezvousBackend

Store rendezvous state behind compare-and-set operations.

tensorplay.distributed.elastic.rendezvous.create_handler

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.

tensorplay.distributed.elastic.rendezvous.RendezvousError

Base class for all rendezvous failures.

tensorplay.distributed.elastic.rendezvous.RendezvousClosedError

The rendezvous has been closed and accepts no more participants.

tensorplay.distributed.elastic.rendezvous.RendezvousTimeoutError

A rendezvous phase exceeded its time budget.

tensorplay.distributed.elastic.rendezvous.RendezvousConnectionError

The rendezvous backend store could not be reached.

tensorplay.distributed.elastic.rendezvous.RendezvousStateError

The rendezvous state is corrupt or the backend rejected an update.

tensorplay.distributed.elastic.rendezvous.RendezvousGracefulExitError

Raised to unwind an agent when it is the last node leaving.

tensorplay.distributed.elastic.rendezvous.RendezvousExhaustedError

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

tensorplay.distributed.elastic.multiprocessing.start_processes

tensorplay.distributed.elastic.multiprocessing.PContext

Base class owning a homogeneous group of worker processes.

tensorplay.distributed.elastic.multiprocessing.RunProcsResult

Outcome of monitoring a worker group to completion.

tensorplay.distributed.elastic.multiprocessing.ProcessFailure

Structured failure of one worker process.

tensorplay.distributed.elastic.multiprocessing.ChildFailedError

Raised by the launcher when one or more workers failed.

tensorplay.distributed.elastic.multiprocessing.SignalException

tensorplay.distributed.elastic.multiprocessing.Redirects

Per-stream redirection modes for a worker group.

tensorplay.distributed.elastic.multiprocessing.Std

Which standard streams a worker's output should go to.

tensorplay.distributed.elastic.multiprocessing.to_map

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

tensorplay.distributed.elastic.utils.data.CyclingIterator

Wrap a finite iterable into an endless one.

tensorplay.distributed.elastic.utils.data.ElasticDistributedSampler

Sampler for elastic jobs where num_replicas may change per epoch.

tensorplay.distributed.elastic.utils.get_env_variable_or_raise

Return the value of env_name or raise if it is unset/empty.

tensorplay.distributed.elastic.utils.get_socket_with_port

Create a socket bound to an ephemeral local port and return it.

tensorplay.distributed.elastic.utils.macros

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

tensorplay.distributed.elastic.events.Event

A single agent or worker lifecycle event.

tensorplay.distributed.elastic.events.NodeState

Lifecycle state carried by events.

tensorplay.distributed.elastic.events.RdzvEvent

A rendezvous progress event, rendered as JSON by the logging handler.

tensorplay.distributed.elastic.events.record

Dispatch event to the handler registered under destination.

tensorplay.distributed.elastic.metrics.configure

Route group (or all groups) to handler.

tensorplay.distributed.elastic.metrics.MetricsConfig

Metric plumbing configuration; maps groups to handler names.

tensorplay.distributed.elastic.metrics.prof

Decorator measuring wall time of fn and emitting it as a metric.

tensorplay.distributed.elastic.timer.TimerClient

Client side of the timer contract.

tensorplay.distributed.elastic.timer.TimerServer

Watches outstanding deadlines and reacts when they expire.

tensorplay.distributed.elastic.timer.configure

Set the process-wide default TimerClient.

tensorplay.distributed.elastic.timer.expires

Context manager asserting the block finishes within after seconds.

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.

On this page

Ask DeepWiki