The workload
Suppose for each user I need to run 20 predictions, each prediction takes 1 second, and there are 150M users. That is a per-item, per-user batch scoring job — the shape of almost every nightly model-refresh, embedding rebuild, or offline ranking pass I have shipped. Before the arithmetic, I separate two quantities people conflate:- Latency is the wall-clock time for one unit of work: here,
20 predictions × 1 s = 20 s per user. - Throughput is units per time: a single worker doing 1 s/prediction sustains
1 prediction/s, or0.05 users/s.
Baseline: sequential, single-threaded
Running the pipeline sequentially on a single thread takes 3,000,000,000 seconds, roughly 95.06 years. The step-by-step:- Total predictions:
150,000,000 users × 20 predictions = 3,000,000,000 predictions. - Total time in seconds:
3,000,000,000 × 1 s = 3,000,000,000 s. - Conversion to practical units:
- Minutes:
3,000,000,000 ÷ 60 = 50,000,000 minutes - Hours:
50,000,000 ÷ 60 ≈ 833,333.33 hours - Days:
833,333.33 ÷ 24 ≈ 34,722.22 days - Years:
34,722.22 ÷ 365.25 ≈ 95.06 years
- Minutes:
Parallel execution runtime
To finish 3 billion predictions in a feasible window, the workload must spread across parallel workers or inference instances (Ray, Spark, or a distributed worker pool). The scaling model I start from is deliberately naive: near-linear scaling with negligible network overhead, soHours ≈ 833,333.33 / workers.
The 50,000-worker row is just
16.67 hrs ÷ 24. If that whole table felt too tidy, good — that is the point of the next section.
Where the naive model lies: bottleneck and serialization
Near-linear scaling assumes every worker contributes one full prediction-second per second and that the work splits perfectly. Real pipelines violate both assumptions, so I model them explicitly. Amdahl’s law bounds the speedup from a fixed serial core. If a fractionp of the job parallelizes and 1-p stays serial, then with N workers the best possible speedup is 1 / ((1-p) + p/N). For example, with 98% parallelizable work (p = 0.98) and N = 10,000:
speedup = 1 / (0.02 + 0.00098) ≈ 48x, not 10,000x.
That serial residue is real: model load/warmup, the driver process writing results, a shared feature store that saturates, or a single coordinator that assigns shards. In the 3B-prediction job above, if 2% is serial, you cannot beat ~48x no matter how many pods you add. I therefore never plan around “number of cores”; I plan around “parallelizable fraction × device throughput.”
The other assumption that breaks is the 1 second per prediction. On a CPU that is often the batch time for one item; group the 20 predictions for a user into a single tensor and the per-item cost collapses. That is the biggest lever, so I treat parallel workers as the last multiplier, not the first.
The three levers, and a combined worked example
The source rules of thumb for shrinking runtime, with the arithmetic I attach to each:- Batching predictions. Running 20 predictions per user one at a time adds launch and Python-loop overhead per item. Grouping the 20 into a single tensor batch drops per-item latency from 1,000 ms to tens of milliseconds on modern hardware. Say 50 ms per batched user.
- Hardware acceleration (GPU/TPU). A model needing 1 second on a CPU core often runs in 5-20 ms on an inference-optimized GPU (NVIDIA L4 or T4) via TensorRT, ONNX Runtime, or vLLM. Say 10 ms per prediction.
- Pre-filtering / pruning. Filtering dormant users and caching static predictions removes work before it enters the pipeline. If half the 150M users can be served from a cached or cheap score, you start from 75M users, not 150M.
20 × 10 ms = 200 ms (I am being conservative and not double-counting the batching speedup):
75M users × 0.2 s = 15,000,000 s of work on one device = ~173.6 days single-stream. Now divide by 10,000 workers: 15,000,000 / 10,000 = 1,500 s ≈ 25 minutes. From 95 years to under half an hour, and the reduction came mostly from device choice and filtering, with parallelism finishing the job — which is the correct order of operations.
What changes: sequential vs parallel stages
So far the whole pipeline is one stage (predict) repeated. Real pipelines are multi-stage: fetch features, transform, embed, score, write. Two regimes:- Sequential stages, one user at a time:
T_user = T_fetch + T_transform + T_embed + T_score + T_write. End-to-end latency is the sum, and throughput is one user per that sum. This is the 95-year shape. - Pipelined / parallel stages: stages overlap via queues, so throughput is set by the slowest stage, not the sum:
T_throughput ≈ max(stage_times). The pipeline that finishes a user everymax(...)seconds is bounded by its bottleneck stage. If scoring is 200 ms but feature fetch is 500 ms/user, adding scoring GPUs does nothing until you fix fetch.
Failure modes I design against
- Stragglers. One slow worker stretches the whole batch’s tail. Shard by key, keep shards equal, and re-launch dropped shards rather than waiting.
- Memory blowup from over-batching. Batch size trades throughput against peak memory; a 150M-user embedding rebuild that batches “as big as possible” OOMs on the first node. Cap batch size by device VRAM, not by optimism.
- Cold-start domination. If model load is 30 s and each worker only scores a few thousand users before the job ends, load time dominates and Amdahl bites hard. Keep workers warm or give each worker enough per-shard work to amortize warmup.
- Idempotency on retry. A distributed pipeline that partially wrote 75M scores and crashed must be resumable; write per-shard checkpoints so a retry does not re-run finished work.