Skip to main content
The single most useful back-of-envelope question in ML platform work is: “will this job ever finish?” I estimate that before I write a line of orchestration code, because a pipeline that takes 95 years on one machine has a fundamentally different design than one that takes two hours. Here is the method I use, worked through a concrete pipeline.

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, or 0.05 users/s.
Total runtime is what falls out of these once you multiply by the population and divide by parallelism.

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:
  1. Total predictions: 150,000,000 users × 20 predictions = 3,000,000,000 predictions.
  2. Total time in seconds: 3,000,000,000 × 1 s = 3,000,000,000 s.
  3. 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
I keep 365.25 (not 365) in the years conversion because leap years over a 95-year span are exactly the kind of thing you want your unit test to catch. The lesson of “95 years” is not that the job is impossible; it is that the sequential execution model is impossible. Every real design adds parallelism, batching, or a faster device — usually all three.

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, so Hours ≈ 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 fraction p 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.
Combine them and the 95-year job stops being science fiction. Take 20 predictions per user, batched on a GPU at 10 ms/prediction so a per-user batch is 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 every max(...) 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.
This is the reason I always identify the bottleneck stage before optimizing anything: in a pipelined system, speedup past the bottleneck is wasted money. It is also why a pipeline can have terrible latency (sum of stages) but great throughput (max of stages) — the two are different problems with different fixes. Parallelizing across users (scale-out workers) and pipelining across stages (producer/consumer) compose multiplicatively, which is exactly the difference between “10,000 workers each doing all stages” and a streaming DAG.
Near-linear scaling in the table ignores three real costs: cold-start/model-load per worker, per-request network and serialization overhead, and shared-resource contention (a feature store or object store that caps aggregate read bandwidth). At 35,000-50,000 workers the coordinator and the data source, not the CPU, become your bottleneck. Budget for the tail.

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.
The same 150M population and the “20 predictions per user” shape drive the serving-fleet and profiling-storage math on User stats; the per-item vs batched latency argument there and here is identical, just viewed from the request path instead of the batch path.