Calculating time to run a pipeline
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, or0.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:
- 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:
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.
| Parallel workers / concurrency | Total hours | Total days |
|---|---|---|
| 100 | 8,333.33 hrs | ~347.2 days |
| 500 | 1,666.67 hrs | ~69.4 days |
| 1,000 | 833.33 hrs | ~34.7 days |
| 5,000 | 166.67 hrs | ~6.9 days |
| 10,000 | 83.33 hrs | ~3.5 days |
| 35,000 | ~23.8 hrs | ~1.0 day |
| 50,000 | 16.67 hrs | ~0.69 days |
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.
Runtime = (predictions_per_run × ms_per_prediction) / parallel_workers
= (150,000,000 × 20 × 1,000 ms) / 1 -> 95.06 years (naive CPU, serial)
= ( 75,000,000 × 20 × 10 ms) / 10,000 -> ~25 minutes (pruned, GPU, parallel)
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.
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.