skip to content

An email-digest sender's ThreadPoolExecutor backs up at a 1,200-send-per-minute peak; why does raising max_workers not fix it?

level: seniorimportance: should knowfreq 45%

answer

  1. Workers are not always the bottleneck
  2. Arrival rate times service time
  3. What blocks the producer when the consumer lags?
  4. The executor's work queue has no maximum
  5. Measure queue wait, not worker count

basics

~20 s

Worker count is not the bottleneck. Threads added past the point where the downstream service saturates only wait somewhere else, while the executor's unbounded work queue hides the growing backlog. Size from arrival rate times service time, and bound admission instead.

solid answer

~50 s

At 1,200 sends a minute that is 20 arrivals a second; if one send takes 200 ms, roughly four sends need to be in flight to keep up, and a queue that grows anyway means service time has risen, not that workers are scarce. Raising `max_workers` past the concurrency the remote endpoint or the shared connection pool will actually serve just moves the wait from the executor's queue into the socket layer. The deeper problem is that the queue behind a `ThreadPoolExecutor` is unbounded, so `submit()` never blocks: the sender happily accepts work it cannot perform, memory grows, and every queued digest ages before it is sent. The fix is admission control — a `threading.Semaphore` acquired before `submit()`, or a bounded `queue.Queue` and your own feeder — plus measuring queue wait time rather than worker count.

code

python · 18 lines
python
import threading
from concurrent.futures import ThreadPoolExecutor

IN_FLIGHT = threading.Semaphore(8)   # admission limit, not worker count

def send(address):
    return address.encode("utf-8")

def guarded(address):
    try:
        return send(address)
    finally:
        IN_FLIGHT.release()

with ThreadPoolExecutor(max_workers=4) as pool:
    for i in range(1000):
        IN_FLIGHT.acquire()          # blocks the producer once 8 are outstanding
        pool.submit(guarded, f"user{i}@example.com")

go deeper

for a junior

Understand that a pool has two separate numbers: how many items run at once and how many are waiting. Know that submitting more work than the pool can perform does not fail loudly; the extra work simply queues.

for a middle

Be able to derive a worker count from arrival rate times service time rather than from cores, and explain why the executor's queue being unbounded means submit() never pushes back on the producer.

for a senior

Show the diagnosis: separate queue wait from service time in metrics, identify the shared downstream limit, and add admission control with a semaphore or a bounded queue instead of turning the worker count up until memory runs out.

for a principal

Frame pool width as a contract with the dependency. Decide what the service does when it cannot keep up — shed, degrade, or slow the intake — and make that behaviour explicit rather than letting an unbounded queue choose it silently.

### The number that actually sets pool width For an I/O-bound pool, worker count is not a property of the machine at all — it is a property of the traffic. Little's law gives it directly: the concurrency you need is arrival rate multiplied by service time. A digest sender at a 1,200-send-per-minute peak is taking 20 arrivals per second; if a send completes in 200 ms, then 20 × 0.2 = 4 sends need to be in flight on average. Twenty workers do nothing for that load except sit idle. If service time triples to 600 ms, the requirement rises to 12, and a pool of 4 begins to fall behind — the same pool, the same code, a different answer, because the input changed. This is why `os.process_cpu_count()` is the wrong input here. Cores bound how much *computation* can happen at once; they say nothing about how many sockets can be open waiting for a remote reply. A pool sized from cores on an I/O workload is either wildly too small on a one-CPU allocation or wildly too large on a wide host. ### Why adding workers stops helping Every I/O pool sits in front of something with its own ceiling: a connection pool with a fixed size, an endpoint enforcing a rate limit, a name resolver, a disk. Once that ceiling is reached, extra workers do not increase completed work per second; they increase the number of places work is waiting. Latency per item gets worse because each in-flight request now shares the same bottleneck with more competitors, and error rates often get worse too, because a saturated remote starts refusing or timing out. The classic shape is a throughput curve that flattens and then bends downward as the pool grows — and past that point, worker count is a knob that only adds memory and scheduler work. Adding threads has a floor cost as well: each is an OS thread with a real stack reservation, and for pure-Python work they contend for the GIL. But the decisive fact is that once the downstream saturates, the pool's width has stopped being the limit. ### The unbounded queue is the real hazard The more damaging half of this failure is invisible from the worker count. The work queue inside a `ThreadPoolExecutor` is unbounded. `submit()` therefore *never* blocks: it appends and returns a `Future` immediately. A sender that loops over a recipient list submits all of them in a fraction of a second, and the queue absorbs every one — each holding a reference to its arguments, and each `Future` holding its eventual result alive. Three consequences follow. First, memory grows with the backlog rather than with the concurrency, and a burst that lasts a few minutes can be the thing that ends the process. Second, latency measured at the send call is meaningless: nearly all the elapsed time is queue wait, which no timing around the actual I/O will show. Third, and most subtly, the system has lost backpressure — the producer never learns that the consumer is behind, so it cannot slow the intake, shed low-priority digests, or fail fast. A backlog of queued work also ages badly: if a batch is rejected downstream because of an encoding mismatch in the payload and each item retries, the retries land behind an already deep queue, and the effective delay per item compounds. ### What to do instead Bound admission. The smallest correct change is a counting semaphore sized to the number of items you are willing to have outstanding: acquire it on the submitting thread before `submit()` and release it in the task's `finally`. Because `acquire()` blocks, the producer now inherits the consumer's rate, which is exactly what backpressure means. The alternative shape is your own bounded `queue.Queue` with a fixed set of consumer threads, where `put()` blocks when full; that also gives you a natural place to drop or prioritise when the queue is at its limit, which a semaphore does not. Then measure the right things. Worker count and CPU utilisation both look healthy while a pool drowns. The two numbers that tell the truth are how long an item waits between being submitted and starting, and how deep the backlog is; a rising queue-wait time with idle workers means the downstream is the limit, while a rising queue-wait time with all workers busy and the downstream unsaturated means the pool is genuinely too narrow. Only in that second case does raising the worker count help — and even then, only up to the concurrency the dependency has agreed to serve.

  • How would you tell from metrics whether the pool is too narrow or the dependency is saturated?
    Record how long each item waits between submission and start, and how many workers are busy. Rising queue wait with workers frequently idle means they are blocked on something shared downstream, so widening the pool changes nothing. Rising queue wait with every worker busy and the dependency still responding within its normal latency means the pool really is too narrow. Worker count and machine CPU look fine in both cases, which is why neither is a useful signal.
  • What is the practical difference between a semaphore around submit() and a bounded queue.Queue with your own consumer threads?
    Both block the producer, which is the point. The semaphore is a two-line change that keeps the executor and its Future objects. A bounded queue.Queue gives you an explicit place to decide what happens when the limit is reached: put() can block, time out, or fail so the caller can drop or downgrade an item. Choose the semaphore for simple throttling and the bounded queue when you need shedding, prioritisation, or a queue depth you can report.
  • Why does memory grow with the backlog even though the pool width is fixed?
    Because the queue, not the pool, holds the backlog. Each queued work item keeps its callable and arguments alive, and each Future stays referenced until the caller drops it, holding any result or exception with it. Width bounds only how much runs at once. A burst of submissions therefore turns into resident memory proportional to the burst, which is precisely what bounding admission prevents.

Adding workers to a saturated dependency is like opening more checkout lanes onto one exit door: the line moves off the tills and forms somewhere you are not looking.

saying these in an interview costs you the question

  • Answers every backlog with a larger max_workers
  • Believes submit() blocks when all workers are busy
  • Sizes an I/O-bound pool from the CPU count
  • Ignores queue wait time when measuring latency
  • Thinks an unbounded queue is harmless because tasks are small
  • Never considers the downstream service's own concurrency limit

context