TensorPlay
Reference guides
latest (dev)
Copy
View Markdown

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 a Future back immediately; the result is produced later.

  • remote() — run the function on the target and get an RRef handle 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)() calls func on the worker to and returns its result. to is a worker name (or WorkerInfo).

  • rpc_async() does the same but returns a Future. 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 the RRef that comes back points at a value that lives on that worker. An RRef can 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 alias new_method()) creates a bound callable to a remote method so you can write obj.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

On this page

Ask DeepWiki