skip to content

How does TensorFlow's MultiWorkerMirroredStrategy find its peer workers?

level: middleimportance: should knowfreq 52%

answer

  1. An environment variable, not a launcher
  2. JSON with two keys
  3. Same script everywhere, different index
  4. Index zero has extra duties
  5. Set it before the strategy is constructed

basics

~20 s

Through the TF_CONFIG environment variable. It is a JSON string holding a "cluster" map of every worker's host:port and a "task" entry giving this process its own type and index. Every worker runs the same script with a different index.

solid answer

~50 s

`MultiWorkerMirroredStrategy` reads the `TF_CONFIG` environment variable, a JSON string with two keys: `cluster`, listing the addresses of every task in the job (typically `{"worker": ["host1:port", "host2:port"]}`), and `task`, identifying *this* process — `{"type": "worker", "index": 0}`. Every worker launches the **same program** with an identical `cluster` and a different `index`; there is no launcher process that hands out ranks. `TF_CONFIG` must be set before the strategy is constructed, because the constructor parses it and immediately starts connecting to the peers — it blocks until all workers have come up, so a job with a typo'd address simply hangs rather than erroring. Worker index 0 is the chief by default and takes on side duties such as writing checkpoints and TensorBoard logs. Because the strategy is synchronous, all workers must run the same number of steps; if one dies, the collective stalls and the job must be restarted.

code

python · 15 lines
python
import json
import os

# Must be set before the strategy is constructed.
os.environ["TF_CONFIG"] = json.dumps(
    {
        "cluster": {"worker": ["10.0.0.1:12345", "10.0.0.2:12345"]},
        "task": {"type": "worker", "index": 0},
    }
)

import tensorflow as tf

strategy = tf.distribute.MultiWorkerMirroredStrategy()
print(strategy.num_replicas_in_sync)

go deeper

for a junior

Be able to say that multi-worker TensorFlow is configured through the TF_CONFIG environment variable, and that every machine runs the same script with a different task index.

for a middle

Describe the JSON shape — cluster addresses plus this task's type and index — explain that it must be set before the strategy is constructed, and note that worker 0 acts as chief.

for a senior

Show operational experience: diagnosing a blocking peer handshake as a hang rather than an error, the temp-directory checkpoint discipline for non-chief workers, and restart-based recovery with BackupAndRestore.

for a principal

Own the surrounding decisions — gang scheduling and its idle-capacity cost, whether a synchronous multi-worker job suits a preemptible fleet, and how much interconnect the job needs before adding workers stops paying.

## The configuration surface Multi-worker TensorFlow has no separate launcher binary that assigns ranks. Instead each process learns its identity from an environment variable, `TF_CONFIG`, whose value is a JSON object with two keys: - **`cluster`** — a map from task type to a list of `host:port` addresses. For `MultiWorkerMirroredStrategy` this is normally just `{"worker": [...]}`. The list must be **identical on every worker**, including its order, because position in the list *is* the worker's identity. - **`task`** — `{"type": "worker", "index": N}`, where `N` is this process's position in that list. So a two-machine job is the same script started twice with the same `cluster` and `index` 0 on one host and 1 on the other. Whatever orchestrates the job — a scheduler, a shell script, a Kubernetes job template — is responsible for injecting the right value; TensorFlow only reads it. `tf.distribute.cluster_resolver.TFConfigClusterResolver` is the object that parses it, and you can pass a resolver explicitly if your cluster information comes from somewhere else. ## Set it before the strategy exists The strategy constructor parses `TF_CONFIG`, opens the gRPC server for this task, and begins the peer handshake. Setting the variable after `MultiWorkerMirroredStrategy()` has been constructed does nothing. In notebooks, where the variable is often set in a cell after TensorFlow has already been imported and used, this is a frequent source of "it silently ran single-worker" confusion — check `strategy.num_replicas_in_sync` if you are unsure whether the peers were found. ## The handshake blocks Construction does not return until every worker in `cluster` has been reached. That is a deliberate design: collectives require all participants. The operational consequence is that a mistake in a hostname, a firewalled port, or one worker pod that has not been scheduled yet does not produce an error — it produces a hang, on every worker. When debugging a stuck multi-worker start, the first questions are always: is every address in the list resolvable and reachable from every other worker, and has every task actually been started? ## What the chief does By convention the task with `index` 0 in the `worker` list is the **chief**. It is a full training worker like all the others, plus it owns cluster-wide side effects: writing the authoritative checkpoint, writing TensorBoard summaries. This matters because writing from several workers to the same path corrupts it. The TensorFlow multi-worker guidance is that *every* worker must call the save API — because saving involves collective ops and skipping it on non-chiefs would deadlock — but non-chief workers must write to a **unique temporary directory** and delete it afterwards, leaving only the chief's output. Note that a separate `"chief"` task type also exists in `TF_CONFIG` and is used by some other strategies; with `MultiWorkerMirroredStrategy` the usual arrangement is worker 0. ## Collective communication Workers exchange gradients with collective ops. You can hint at the implementation with `tf.distribute.experimental.CommunicationOptions(implementation=tf.distribute.experimental.CommunicationImplementation.NCCL)` — the alternatives are `RING` and `AUTO`, the default. NCCL generally wins on GPU hosts with fast interconnect; `RING` is the fallback where NCCL is unavailable or unreliable. This is a tuning knob, not a correctness one. ## All-or-nothing execution Synchronous training means every worker must execute the same number of steps: each step includes an all-reduce that no worker can complete alone. Consequences worth stating in an interview: - **Gang scheduling.** The job needs all its workers simultaneously; N-1 workers do no useful work. - **The slowest worker sets the pace.** A heterogeneous fleet wastes the fast machines. - **A worker crash stops everything.** There is no in-place recovery. For the last point, the mitigation is `tf.keras.callbacks.BackupAndRestore`, which snapshots training state to `backup_dir` and, when the whole job is restarted, resumes from the beginning of the interrupted epoch rather than from scratch. It is restart-based recovery layered on an external supervisor, not live fault tolerance. ## Data and batching The global batch size now spans every replica on every worker, so a two-worker job with four GPUs each has 8 replicas. Each worker builds its own input pipeline and receives a shard of the data, which is the subject of auto-sharding.

  • A multi-worker job starts and then hangs with no error. What is your first hypothesis?
    A peer was never reached. The strategy constructor blocks until every address in the `cluster` list responds, so an unreachable host, a blocked port, or a worker task that was never started produces a hang rather than an exception. Check that every address resolves from every worker and that all tasks were actually launched.
  • How should checkpoints be written in a MultiWorkerMirroredStrategy job?
    Every worker must call the save, because saving involves collective ops and skipping it on non-chiefs deadlocks the job. Only the chief writes to the real destination; the other workers write to unique temporary directories and delete them afterwards. Writing several workers to one path corrupts the checkpoint.
  • What happens if one worker of eight crashes mid-epoch?
    The remaining workers block on the next all-reduce and the job stalls — synchronous collectives need every participant. There is no in-place recovery, so the supervisor restarts the whole job. `tf.keras.callbacks.BackupAndRestore` makes that restart cheap by resuming from the interrupted epoch instead of from scratch.
  • Can you control which collective implementation the workers use?
    Yes — pass `tf.distribute.experimental.CommunicationOptions` with `implementation` set to `CommunicationImplementation.NCCL`, `RING`, or the default `AUTO`. NCCL usually wins on GPU hosts with a fast interconnect; RING is the portable fallback. It is a performance knob, not a correctness one.

saying these in an interview costs you the question

  • Expecting a launcher to assign worker ranks
  • Setting TF_CONFIG after building the strategy
  • Giving workers different cluster lists or orders
  • Having only the chief call the checkpoint save
  • Believing a dead worker is routed around automatically

context