> ## Documentation Index
> Fetch the complete documentation index at: https://authorsnote.askailab.online/llms.txt
> Use this file to discover all available pages before exploring further.

# Calculating time to run a pipeline

> How I estimate end-to-end runtime and throughput for a multi-stage data/AI pipeline, and what changes when stages run in parallel.

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`.

| 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.

```text theme={null}
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 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.

<Warning>
  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.
</Warning>

## 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](/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.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.