Worker Pools with Priority Lanes¶
A worker pool that serves both urgent jobs and a batch backlog needs lanes, and the way workers are assigned to lanes decides what happens when either side gets busy. Measured on Python 3.14 with a simulated pool — 10 workers, batch jobs taking 100–300 ms, urgent jobs taking 20 ms, and a backlog of 2,000 batch jobs queued at the start: with one shared FIFO queue, none of 240 urgent jobs finished within the 8-second run. Strict priority — workers always take an urgent job first — served all of them, but urgent p99 was 104 ms, because a job must wait for a worker to finish its current batch item. Two workers dedicated to urgent jobs cut the p99 to 21.4 ms — until urgent traffic rose to 150 per second, beyond what two workers can do, when only 832 of 1,200 finished and p99 reached 2,924 ms. Two reserved workers plus eight shared workers that prefer urgent jobs kept p99 at 21.4 ms under normal load and 54 ms under the spike, at the cost of batch throughput falling from 41 to 35 jobs per second. This guide builds each arrangement and shows how to size the reserve.
Prerequisites¶
- A worker pool, from building an async worker pool with TaskGroup.
- Priority at a shared limit, from prioritizing requests under a shared rate limit.
- The topic overview, Worker Pool Implementations.
1. Give each lane its own queue¶
Start by separating the work. One queue per lane lets workers choose, and lets you measure each lane's depth:
class Lanes:
def __init__(self):
self.urgent: asyncio.Queue = asyncio.Queue()
self.batch: asyncio.Queue = asyncio.Queue()
self.wake = asyncio.Event()
def submit(self, job, urgent: bool):
(self.urgent if urgent else self.batch).put_nowait(job)
self.wake.set()
def take(self, lanes: tuple[str, ...]):
for name in lanes: # lanes in preference order
q = getattr(self, name)
if not q.empty():
return q.get_nowait()
return None
Measured with a single shared FIFO queue instead: the 2,000 batch jobs needed about 40 seconds of the pool's time, and every urgent job queued behind them. In the 8-second run, 0 of 240 urgent jobs at 30 per second completed. A lane is the minimum structure that lets urgent work be seen.
Verify: export the depth of each lane's queue separately; one combined depth hides an urgent backlog behind a batch one.
2. Measure strict priority and its limit¶
The simplest policy has every worker check the urgent lane first:
async def worker(lanes: Lanes, prefer=("urgent", "batch")):
while True:
job = lanes.take(prefer)
if job is None:
lanes.wake.clear()
await lanes.wake.wait()
continue
await job.run()
Measured at 30 urgent jobs per second: all 240 completed, with a median of 32.4 ms and a p99 of 104 ms from submission to completion. The job itself takes 20 ms; the rest is waiting. The pool does not preempt: an urgent job that arrives while all 10 workers are inside batch jobs waits until one finishes, and with batch jobs of 100–300 ms, that wait reached about 80 ms at p99. Batch throughput was 48.1 per second. Strict priority is enough when batch jobs are short compared with the urgent latency target, and not otherwise.
Verify: compare urgent p99 with the longest batch job divided by the worker count; if they are close, non-preemption is setting the tail.
3. Reserve workers for the urgent lane¶
To remove the wait for a batch job to finish, keep some workers that never take batch work. They are idle when no urgent work is waiting, which is their purpose:
async def run_pool(lanes: Lanes, workers=10, reserved=2):
async with asyncio.TaskGroup() as tg:
for _ in range(reserved):
tg.create_task(worker(lanes, prefer=("urgent",)))
for _ in range(workers - reserved):
tg.create_task(worker(lanes, prefer=("batch",))) # dedicated lanes
Measured at 30 urgent jobs per second: median 20.2 ms and p99 21.4 ms — the job's own 20 ms plus scheduling — and batch throughput fell to 41.5 per second, since two workers no longer did batch work. Then the spike: at 150 urgent jobs per second, two workers can complete at most 100, and the urgent lane's queue grew without bound. Only 832 of 1,200 urgent jobs finished in the run, with a median of 1,487 ms and a p99 of 2,924 ms, while eight batch-only workers carried on at 42.5 per second. Fully separated lanes isolate the batch from urgent spikes — and the urgent lane from any help.
Verify: the reserved workers' capacity, reserved / urgent job time, is above peak urgent arrival rate — or the shared workers can help, as in the next step.
4. Let shared workers help the urgent lane¶
The arrangement that held under both loads combines the two: reserved workers take only urgent jobs, and the remaining workers prefer urgent jobs and take batch work otherwise:
async def run_pool(lanes: Lanes, workers=10, reserved=2):
async with asyncio.TaskGroup() as tg:
for _ in range(reserved):
tg.create_task(worker(lanes, prefer=("urgent",)))
for _ in range(workers - reserved):
tg.create_task(worker(lanes, prefer=("urgent", "batch")))
Measured at 30 per second: identical to dedicated lanes — 20.2 ms median, 21.4 ms p99, 41.2 batch jobs per second — because the reserved pair handled all the urgent work. At 150 per second: all 1,200 urgent jobs finished, with a median of 22.5 ms and a p99 of 54.0 ms, as shared workers picked up urgent jobs whenever they finished a batch job. Batch throughput fell to 35.4 per second during the spike — the cost of absorbing it — and recovered when it ended.
Verify: in a load test with urgent traffic above the reserved capacity, urgent jobs all complete and batch throughput drops but does not stop.
5. Size the reserve from arrival rates¶
The reserve should cover normal urgent load, so that urgent latency does not depend on batch job length; spikes are what the shared workers are for. Compute it from measured rates and add headroom:
def reserved_workers(urgent_rate_per_s: float, urgent_job_s: float, headroom=1.5) -> int:
return max(1, math.ceil(urgent_rate_per_s * urgent_job_s * headroom))
# 30/s x 20 ms = 0.6 busy workers -> 1 reserved; the 2 used above leave room for bursts
Export, per lane, queue depth, wait time from submission to start, and completions per second. Urgent wait rising above the job time means the reserve is too small for current load; batch completions dropping to zero for long periods means urgent traffic is consuming the whole pool, and the pool needs more workers rather than a different split. For more than two classes — per-tenant fairness, for example — the same structure extends to weighted lanes, as in fair scheduling across tenants in a worker pool.
Verify: reserved capacity exceeds normal urgent arrival rate times job time, with headroom, and per-lane metrics are on a dashboard.
Verification¶
A pool with priority lanes is working when:
- Each lane has its own queue and depth metric.
- Urgent p99 is close to the urgent job's own duration — 21.4 ms for 20 ms jobs here — not to batch job length.
- Urgent spikes above reserved capacity still complete, because shared workers help.
- Batch throughput degrades during spikes but recovers afterwards.
Diagnostic Hook: when urgent latency follows the batch job duration — p99 of 104 ms with 100–300 ms batch jobs here — urgent jobs are waiting for workers to finish batch items. Reserve workers for the urgent lane; strict priority alone cannot preempt a running job.
Pitfalls & edge cases¶
- One FIFO queue. Measured: 0 of 240 urgent jobs finished behind a 2,000-job backlog.
- Strict priority with long batch jobs. Measured: urgent p99 104 ms for a 20 ms job.
- Dedicated lanes only. Measured: a 150/s spike left 368 urgent jobs unfinished.
- Reserving too many workers. Each one removes batch capacity, 41 against 48 jobs per second here.
Frequently Asked Questions¶
How do I prioritize jobs in an asyncio worker pool?
Use one queue per lane, reserve a few workers that only take urgent jobs, and let the remaining workers take urgent jobs first and batch jobs otherwise.
Why is strict priority not enough for urgent jobs?
Workers do not preempt running jobs, so an urgent job waits for one to finish. With 100-300 ms batch jobs, urgent p99 was 104 ms for a 20 ms job.
What happens when urgent traffic exceeds dedicated workers?
Its queue grows without bound. Two dedicated workers handle 100 jobs/s; at 150/s only 832 of 1,200 finished, with p99 2.9 s.
How many workers should be reserved for urgent jobs?
Urgent arrival rate times job duration, plus headroom: 30/s × 20 ms needs 0.6 workers, so one or two. Shared workers absorb spikes above that.
Related¶
- Worker Pool Implementations — up to the topic overview.
- Rate-limited worker pools — when the lanes share an upstream limit.
- Concurrent Execution & Worker Patterns — the section overview.