Skip to content

Scaling a Slow Pipeline Stage

A staged pipeline runs at the speed of its slowest stage, and adding workers anywhere else changes nothing. The work is finding that stage, widening it by the right amount, and seeing where the bottleneck moves next. Measured on Python 3.14 with a three-stage pipeline — fetch (5 ms of I/O), parse (1.2 ms of pure-Python CPU) and store (10 ms of I/O) — over 2,000 items: with one worker per stage it ran at 95 items/s, and the store stage was busy 100% of the time. Quadrupling the fetch workers left it at 95/s. Four store workers gave 376/s, sixteen gave 692/s — at which point fetch was the stage at 100%. With eight fetchers and four parse workers on the event loop it reached only 816/s: the four parse workers shared one thread, each busy 25%, together saturating it, and loop lag rose to 9.8 ms. Moving parse to a four-process pool gave 1,412/s with loop lag at 1.0 ms, and doubling every stage gave 2,793/s — 29 times the starting point. This guide is the procedure behind those steps.

Prerequisites

1. Measure each stage's utilization

Queue depth alone is ambiguous: in a pipeline with bounded queues, every queue upstream of the bottleneck fills up, so several queues look full at once. Measuring how busy each stage's workers are is unambiguous. Accumulate the time each worker spends processing — not waiting on queues — and divide by workers × elapsed time:

busy = {"fetch": 0.0, "parse": 0.0, "store": 0.0}

async def store_worker(q_in):
    while (record := await q_in.get()) is not None:
        start = time.perf_counter()
        await store(record)                         # the work, not the queue wait
        busy["store"] += time.perf_counter() - start

def utilization(elapsed: float, workers: dict[str, int]) -> dict[str, float]:
    return {stage: busy[stage] / (workers[stage] * elapsed) for stage in busy}

Measured with one worker per stage: fetch 48%, parse 13%, store 100%, and all three input queues at their 50-item capacity. The store stage was the bottleneck; the full queues ahead of it were its backlog. The throughput agrees: 95 items/s is one store worker's limit at about 10.5 ms per item including overhead. Adding fetch workers in that state — the intuitive move, because the first queue is full — was measured too: four fetchers, still 95/s, fetch utilization down to 12%.

Verify: your pipeline reports utilization per stage, and one stage is near 100% while the throughput matches that stage's capacity.

Throughput as each bottleneck was widened 7 horizontal bars comparing 1/1/1 with the others. Throughput as each bottleneck was widened 1/1/1 95/s 4 fetch 95/s 4 store 376/s 16 store 692/s 8 fetch, 4 parse on loop 816/s parse in 4 processes 1,412/s 16/8/32, 8 processes 2,793/s fetch 5 ms I/O, parse 1.2 ms CPU, store 10 ms I/O; 2,000 items; Python 3.14. Only widening the stage at 100% raised throughput.

2. Size the stage with Little's law

For an I/O-bound stage, the number of workers needed follows from the target rate and the time each item takes: workers = rate × time per item. That is Little's law applied to the stage:

def workers_needed(target_per_s: float, seconds_per_item: float, headroom: float = 1.3) -> int:
    return math.ceil(target_per_s * seconds_per_item * headroom)

workers_needed(700, 0.0105)    # store stage: 10 workers for 700 items/s

Measured: four store workers gave 376/s — four times one worker, as the law predicts — and sixteen gave 692/s with store utilization at 49%, meaning store had capacity to spare and was no longer the limit. Fetch, at 4 workers × 5 ms, had a ceiling of 800/s, and was now at 100%. Headroom matters because time per item is not constant; a stage sized to exactly 100% at the average queues up on every slow item. For I/O stages, extra asyncio workers are cheap — a few kilobytes each — so the real limit is what the downstream system tolerates, which is where limiting concurrent requests with asyncio.Semaphore belongs.

Verify: after resizing, the stage's utilization drops below about 80% and throughput rises to the next stage's ceiling.

3. Recognise a CPU-bound stage on the loop

A CPU-bound stage running on the event loop does not behave like an I/O stage: adding workers adds no capacity, because they all share one thread. It also hides from the queue-depth signal, because it slows every other stage at the same time:

async def parse_worker(q_in, q_out):
    while (item := await q_in.get()) is not None:
        await q_out.put(parse_cpu(item))           # 1.2 ms of Python, on the loop

Measured with eight fetchers, four parse workers and sixteen store workers: 816 items/s. Each parse worker reported 25% utilization — which looks like spare capacity until you add them up: four workers × 25% is one thread at 100%. At 816 items/s × 1.2 ms the parse work alone used the whole loop, and the maximum loop lag measured by a sampling task rose to 9.8 ms, the highest of any configuration. Fetch reported 99% because its 5 ms sleeps were being stretched by the busy loop, not because it needed more workers. The tell-tale signs are the sum of a stage's utilization reaching about 100% of one core, and loop lag rising with throughput.

Verify: for each stage, you know whether its work is I/O or CPU, and the summed utilization of CPU stages on the loop stays well below one core.

Utilization per stage at each step A grid of 5 rows by 6 columns. Utilization per stage at each step configuration items/s fetch parse store max loop lag 1 / 1 / 1 95 48% 13% 100% 1.9 ms 4 / 1 / 16 692 100% 85% 49% 4.6 ms 8 / 4 on loop / 16 816 99% 4 x 25% 88% 9.8 ms 8 / 4 processes / 16 1,412 99% 60% 91% 1.0 ms 16 / 8 processes / 32 2,793 97% 82% 91% 1.3 ms Workers per stage are fetch / parse / store; the bottleneck stage is the one at 100%.

4. Move CPU-bound stages off the loop

A CPU stage scales with processes, not coroutines. Keep the stage's workers as coroutines — they handle queueing and backpressure — and have each one hand the CPU work to a process pool sized to match:

from concurrent.futures import ProcessPoolExecutor

PARSE_WORKERS = 4
pool = ProcessPoolExecutor(PARSE_WORKERS)

async def parse_worker(q_in, q_out):
    loop = asyncio.get_running_loop()
    while (item := await q_in.get()) is not None:
        await q_out.put(await loop.run_in_executor(pool, parse_cpu, item))

Measured: 1,412 items/s with four processes, loop lag back to 1.0 ms, parse utilization 60%. Fetch was again the stage at 99%, so the next step widened everything: 16 fetchers, 8 parse processes and 32 store workers gave 2,793 items/s with loop lag at 1.3 ms. One coroutine per process is the right ratio — more coroutines than processes only queue inside the executor. Process pools add pickling cost per item, which matters when items are large; reducing pickle overhead in ProcessPoolExecutor payloads covers how to keep it small.

Verify: after moving the stage, loop lag stays near its idle level while throughput rises.

The scaling loop A flow of 5 stages. The scaling loop Measure utilization per stage, loop lag Find the 100% stage not the fullest queue I/O stage workers = rate x time x 1.3 CPU stage process pool, 1 coroutine each Measure again the bottleneck moved Four rounds took this pipeline from 95 to 2,793 items/s.

5. Stop at the limit that is not yours to raise

Each round moves the bottleneck, and eventually it lands on something outside the pipeline: the database's write capacity, an API's rate limit, the network. Past that point, more workers only add contention. Make the external limit explicit in the stage that touches it, and let the pipeline's bounded queues turn the excess into backpressure:

STORE_LIMIT = asyncio.Semaphore(32)       # what the database team agreed to

async def store_worker(q_in):
    while (record := await q_in.get()) is not None:
        async with STORE_LIMIT:
            await store(record)

When the store stage is capped and at 100%, the right next step is not more store workers but fewer round trips per item — writing in batches, as in batching database writes in async pipelines. Keep the utilization metrics in production: input data changes, item sizes grow, and the stage that was the bottleneck at launch is rarely the one six months later.

Verify: every stage that calls an external system has an explicit concurrency cap, and dashboards show utilization per stage.

How do I speed up this stage? A decision on What is true of the stage with 4 outcomes. How do I speed up this stage? What is true of the stage? below ~80% busy leave it alone 4 fetchers: still 95/s I/O, at 100% more workers, Little's law 4 to 16 store: 376 to 692/s CPU, at 100% of a core process pool 816 to 1,412/s at an external limit batch, fewer round trips more workers add contention Widen only the stage that is saturated.

Verification

A pipeline stage is scaled correctly when:

  • Utilization is measured per stage, excluding queue waits.
  • The widened stage was the one at 100%, and throughput rose to the next ceiling.
  • CPU stages run in processes, one coroutine per process, with loop lag unchanged.
  • External limits are explicit caps, and saturation there leads to batching rather than more workers.

Diagnostic Hook: export busy time per stage as a counter and graph its rate divided by worker count. The stage whose line sits at 1.0 is the bottleneck; when a deploy moves that line from one stage to another, the throughput change that follows is no mystery.

Pitfalls & edge cases

  • Widening the stage with the fullest queue. Measured: four fetchers, still 95/s.
  • Adding coroutines to a CPU stage. Four parse workers at 25% each were one saturated thread.
  • Trusting utilization while the loop is saturated. I/O stages look busier than they are.
  • Scaling past an external limit. Throughput stays put and the downstream suffers.

Frequently Asked Questions

How do I find the bottleneck in an asyncio pipeline?

Measure the fraction of time each stage's workers spend processing, not waiting on queues. The stage at about 100% is the bottleneck; full queues upstream of it are its backlog. In testing, the store stage was at 100% while all three queues were full.

How many workers does a pipeline stage need?

Target rate times time per item, plus headroom: for 700 items/s at 10.5 ms each, about 10. Four workers at 10 ms gave 376 items/s, four times one worker.

Why doesn't adding workers to a CPU-bound stage help?

Coroutines share one thread. Four parse workers each 25% busy were one core at 100%, and loop lag rose to 9.8 ms. Move the work to a process pool.

What should I do when the bottleneck is the database?

Cap concurrency at what the database tolerates and reduce round trips per item, for example by batching writes, rather than adding workers.