# tensorplay.distributed.rpc Source: https://www.tensorplay.cn/docs/distributed.rpc.html 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()](/docs/generated/tensorplay.distributed.rpc.init_rpc.html#tensorplay.distributed.rpc.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()](/docs/generated/tensorplay.distributed.rpc.rpc_sync.html#tensorplay.distributed.rpc.rpc_sync) — fire the function and block until the result is returned. - [rpc_async()](/docs/generated/tensorplay.distributed.rpc.rpc_async.html#tensorplay.distributed.rpc.rpc_async) — fire the function and get a [Future](/docs/generated/tensorplay.distributed.rpc.Future.html#tensorplay.distributed.rpc.Future) back immediately; the result is produced later. - [remote()](/docs/generated/tensorplay.distributed.rpc.remote.html#tensorplay.distributed.rpc.remote) — run the function on the target and get an [RRef](/docs/generated/tensorplay.distributed.rpc.RRef.html#tensorplay.distributed.rpc.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 | tensorplay.distributed.rpc.init_rpc | | | --- | --- | | tensorplay.distributed.rpc.shutdown | | | tensorplay.distributed.rpc.get_worker_info | | | tensorplay.distributed.rpc.get_rpc_timeout | | | tensorplay.distributed.rpc.is_available | | [init_rpc()](/docs/generated/tensorplay.distributed.rpc.init_rpc.html#tensorplay.distributed.rpc.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()](/docs/generated/tensorplay.distributed.rpc.shutdown.html#tensorplay.distributed.rpc.shutdown) tears the agent down, waiting for in-flight calls first when graceful=True. [get_worker_info()](/docs/generated/tensorplay.distributed.rpc.get_worker_info.html#tensorplay.distributed.rpc.get_worker_info) returns the WorkerInfo for a named worker (or for the current worker when no name is given), and [is_available()](/docs/generated/tensorplay.distributed.rpc.is_available.html#tensorplay.distributed.rpc.is_available) reports whether RPC is compiled into this build. ## Remote calls | tensorplay.distributed.rpc.rpc_sync | | | --- | --- | | tensorplay.distributed.rpc.rpc_async | | | tensorplay.distributed.rpc.remote | | | tensorplay.distributed.rpc.RRef | | | tensorplay.distributed.rpc.Future | | | tensorplay.distributed.rpc.AllGatherStates | | | tensorplay.distributed.rpc.method_factory | | | tensorplay.distributed.rpc.new_method | | - 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()](/docs/generated/tensorplay.distributed.rpc.rpc_async.html#tensorplay.distributed.rpc.rpc_async) does the same but returns a [Future](/docs/generated/tensorplay.distributed.rpc.Future.html#tensorplay.distributed.rpc.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()](/docs/generated/tensorplay.distributed.rpc.remote.html#tensorplay.distributed.rpc.remote) creates a remote reference. The function runs on the target worker and the [RRef](/docs/generated/tensorplay.distributed.rpc.RRef.html#tensorplay.distributed.rpc.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()](/docs/generated/tensorplay.distributed.rpc.method_factory.html#tensorplay.distributed.rpc.method_factory) (and its alias [new_method()](/docs/generated/tensorplay.distributed.rpc.new_method.html#tensorplay.distributed.rpc.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()](/docs/generated/tensorplay.distributed.rpc.rpc_async.html#tensorplay.distributed.rpc.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 | tensorplay.distributed.rpc.BackendType | | | --- | --- | | tensorplay.distributed.rpc.BackendValue | | | tensorplay.distributed.rpc.TensorPipeRpcBackendOptions | | | tensorplay.distributed.rpc.register_backend | | | tensorplay.distributed.rpc.backend_registered | | | tensorplay.distributed.rpc.init_backend | | | tensorplay.distributed.rpc.construct_rpc_backend_options | | [BackendType](/docs/generated/tensorplay.distributed.rpc.BackendType.html#tensorplay.distributed.rpc.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](/docs/generated/tensorplay.distributed.rpc.TensorPipeRpcBackendOptions.html#tensorplay.distributed.rpc.TensorPipeRpcBackendOptions) for TENSORPIPE — that carries settings passed to [init_rpc()](/docs/generated/tensorplay.distributed.rpc.init_rpc.html#tensorplay.distributed.rpc.init_rpc). Registering a custom backend is done with [register_backend()](/docs/generated/tensorplay.distributed.rpc.register_backend.html#tensorplay.distributed.rpc.register_backend), which associates a backend name with the two handlers that build its options and initialize it; [init_backend()](/docs/generated/tensorplay.distributed.rpc.init_backend.html#tensorplay.distributed.rpc.init_backend) then constructs a concrete agent from a backend and a name. [construct_rpc_backend_options()](/docs/generated/tensorplay.distributed.rpc.construct_rpc_backend_options.html#tensorplay.distributed.rpc.construct_rpc_backend_options) builds the default options object for a backend. ## Where to go next - [the distributed package](/docs/distributed.html) — process groups and collectives, the data-parallel counterpart to RPC. - [distributed tensors](/docs/distributed.tensor.html) — the sharded tensor type you can exchange over RPC. - [device mesh](/docs/distributed.device_mesh.html) — defining the rank topology you may want to share with RPC workers.