latest (dev)
Copy
Latest development documentation · Updated 2026-10-08
tensorplay.distributed.rpc
tensorplay.distributed.rpc lets you call functions on a worker running in
another process — possibly on another machine — and get the result back. It
is the primitive behind the higher-level distributed model APIs (remote
module, remote tensors); the source of truth, however, is the worker-to-worker
call.
RPC is set up separately from the process group. You call
init_rpc() on every worker; each worker is
identified by a name, and the workers announce themselves to each other via
the rendezvous store or environment variables so that a call targeting
"worker1" can be routed to the right process. Once initialized you make
three kinds of calls to a target worker:
rpc_sync()— fire the function and block until the result is returned.rpc_async()— fire the function and get aFutureback immediately; the result is produced later.remote()— run the function on the target and get anRRefhandle to its result, which lives on the remote worker and can be passed to further remote calls without copying the value over the wire.
import os
import tensorplay.distributed.rpc as rpc
os.environ["MASTER_ADDR"] = "localhost"
os.environ["MASTER_PORT"] = "29501"
# run on both workers, each with its own name:
rpc.init_rpc("trainer", rank=int(os.environ["RANK"]), world_size=2)
def add(a, b):
return a + b
# on rank 0, ask the other worker to compute
result = rpc.rpc_sync("trainer", add, args=(3, 4))
print(result) # 7
rpc.shutdown()
Initialization and shutdown
init_rpc() initializes the RPC agent: it
takes the worker name, a backend (the only built-in backend is
TENSORPIPE), the worker’s
rank and world_size, and backend-specific options. It blocks until all
workers in the world have joined, then performs an all-gather so every worker
can map names to the others. Calling it a second time, or calling any RPC
function before it, raises. shutdown() tears
the agent down, waiting for in-flight calls first when graceful=True.
get_worker_info() returns the
WorkerInfo for a named worker (or for the
current worker when no name is given), and is_available()
reports whether RPC is compiled into this build.
Remote calls
rpc_sync(to, func, args=None, kwargs=None, timeout=-1)()callsfuncon the workertoand returns its result.tois a worker name (orWorkerInfo).rpc_async()does the same but returns aFuture. The future is the only handle you need: it resolves to the result (or raises the remote exception) when the call completes, and it also carries the stream/tensor asynchrony of the remote worker.remote()creates a remote reference. The function runs on the target worker and theRRefthat comes back points at a value that lives on that worker. AnRRefcan be a call argument, so you can pass a remote object into another remote call; the worker that owns it keeps it alive until the last reference is gone, which is what makes this building the foundation for distributed objects.method_factory()(and its aliasnew_method()) creates a bound callable to a remote method so you can writeobj.method(x)style calls after resolving the method once.
rpc_async() and rpc_sync take a timeout
in seconds; if the target worker does not respond within it the call raises.
Backends and options
BackendType is the enum of built-in RPC
backends; the supported value is TENSORPIPE, a gRPC-based transport. Each
backend has a corresponding options object — TensorPipeRpcBackendOptions
for TENSORPIPE — that carries settings passed to init_rpc().
Registering a custom backend is done with register_backend(),
which associates a backend name with the two handlers that build its options
and initialize it; init_backend() then
constructs a concrete agent from a backend and a name.
construct_rpc_backend_options() builds the
default options object for a backend.
Where to go next
the distributed package — process groups and collectives, the data-parallel counterpart to RPC.
distributed tensors — the sharded tensor type you can exchange over RPC.
device mesh — defining the rank topology you may want to share with RPC workers.
Help improve this page
Found an error, an unclear step, or a missing example?
tensorplay.distributed.pipelining
tensorplay.distributed.pipelining splits one model so its layers run on different workers, one stage of the pipeline per worker, and orchestrates the forward and backward passes as streams of microbatches flowing through
tensorplay.distributed.tensor
Distributed (sharded) tensors built on top of a device mesh . A DTensor is a logical tensor whose data is split across the ranks of a mesh: each rank stores a local shard, and the shards are expected to behave like a sin

