How does TensorFlow's MultiWorkerMirroredStrategy find its peer workers?
answer
- An environment variable, not a launcher
- JSON with two keys
- Same script everywhere, different index
- Index zero has extra duties
- Set it before the strategy is constructed
basics
~20 sThrough 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 linesimport 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
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.
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.
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.
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