Limiting Concurrency per Job Type¶
A worker that runs several kinds of job shares one concurrency budget between them. When a slow, heavy job type floods the queue — month-end reports, a bulk re-index — it occupies every slot, and quick, latency-sensitive jobs like password-reset emails wait behind it. The instinctive fix, wrapping the heavy job in a semaphore, makes it worse. In a test with arq, 30 report jobs of 500 ms and 20 email jobs of 10 ms on one worker with max_jobs=10: with no limit, emails waited a median of 1,545 ms; with a Semaphore(3) inside the report job, 3,706 ms — because a report blocked on the semaphore still holds a worker slot. Making the report job defer itself when the semaphore was busy brought emails to 40 ms; running reports and emails on separate queues gave 36 ms. This guide implements both fixes and explains when each fits.
Prerequisites¶
- Python 3.11+,
pip install arqand Redis; the reasoning applies to Celery and taskiq queues too. - arq workers, from running arq workers with Redis.
- Bulkheads, from bulkhead isolation with per-dependency semaphores.
1. See why a semaphore inside the job backfires¶
report_sem = asyncio.Semaphore(3)
async def build_report(ctx, report_id: int) -> None:
async with report_sem: # at most 3 reports run at once...
await render_report(report_id) # ...but waiting reports still occupy worker slots
The worker pulls jobs in queue order up to max_jobs=10. It pulls ten reports; three run, seven sit inside async with report_sem, each holding a slot. No email can start until reports finish and release slots — and the next jobs pulled are more reports. Measured: email median latency went from 1,545 ms without the semaphore to 3,706 ms with it, and the whole batch took 5.0 s instead of 1.6 s.
A semaphore limits how many jobs do work; it does not limit how many jobs occupy slots. In a job system, slots are the scarce resource.
Verify: count in-progress jobs per type on a loaded worker; if waiting reports hold most of max_jobs, the semaphore is the cause.
2. Fix 1: defer instead of waiting¶
If the limit is reached, give the slot back and come back later. In arq, raise Retry with a short defer:
from arq import Retry
report_sem = asyncio.Semaphore(3)
async def build_report(ctx, report_id: int) -> None:
if report_sem.locked():
raise Retry(defer=0.1) # slot freed now; job re-queued shortly
async with report_sem:
await render_report(report_id)
Each deferred report releases its slot immediately, so emails run as soon as they are reached. Measured: email median 40 ms. Two costs: deferrals count as tries, so set max_tries high enough for the expected backlog (or track deferrals separately from real failures); and the deferred jobs cycle through Redis while they wait, adding queue traffic. The check-then-acquire is safe on a single event loop because there is no await between locked() and async with.
This fix works within one worker process. Across processes, the in-memory semaphore is per worker; a fleet-wide limit needs a shared counter, step 4.
Verify: under the same load, email latency stays near its uncontended value, and reports still never exceed three concurrently per worker.
3. Fix 2: separate queues with separate workers¶
The structural fix is to stop sharing: put each job class on its own queue and give each queue its own workers with their own max_jobs:
class ReportWorker:
functions = [build_report]
queue_name = "q:reports"
max_jobs = 3 # the limit IS the worker's concurrency
class EmailWorker:
functions = [send_email]
queue_name = "q:emails"
max_jobs = 50
# producers choose the queue
await arq.enqueue_job("build_report", report_id, _queue_name="q:reports")
await arq.enqueue_job("send_email", user_id, _queue_name="q:emails")
Measured with 3 report slots and 7 email slots: email median 36 ms. Each class now scales independently — more email workers for a signup spike, more report workers at month end — and a crash or memory blow-up in the report worker cannot touch email delivery. This is the bulkhead pattern applied to job processing, and the default recommendation for job classes with very different latency or resource needs.
Verify: scale report workers to zero; emails keep flowing at full speed.
4. Enforce a fleet-wide limit with Redis¶
Some limits are global: a third-party API that allows five concurrent report exports per account, no matter how many workers you run. Keep the count in Redis with an expiring lease per running job, and defer when the limit is reached:
import uuid
ACQUIRE = """
redis.call('ZREMRANGEBYSCORE', KEYS[1], '-inf', ARGV[1]) -- drop expired leases
if redis.call('ZCARD', KEYS[1]) < tonumber(ARGV[3]) then
redis.call('ZADD', KEYS[1], ARGV[2], ARGV[4])
return 1
end
return 0
"""
async def build_report(ctx, report_id: int) -> None:
r = ctx["redis"]
lease = str(uuid.uuid4())
now = time.time()
ok = await r.eval(ACQUIRE, 1, "limit:reports", now, now + 120, 5, lease)
if not ok:
raise Retry(defer=1.0)
try:
await render_report(report_id)
finally:
await r.zrem("limit:reports", lease)
Each running job holds one member in a sorted set, scored by its lease expiry; a crashed worker's lease ages out after 120 s instead of holding the slot forever. Choose the lease longer than the longest job, or renew it from a heartbeat for long jobs — the technique in renewing lock leases with a heartbeat task.
Verify: run ten workers; the sorted set never holds more than five members, and killing a worker mid-job frees its slot within the lease time.
5. Prioritise without separate infrastructure¶
When the issue is ordering rather than isolation — urgent jobs should jump the line — a priority scheme can be enough. Most job systems support either multiple queues polled in priority order or per-job priorities; with arq, run one worker deployment per queue and size the urgent one generously. In-process, the priority queue from implementing a priority queue with asyncio.Queue does the same for an internal pipeline.
Priority does not stop a flood of low-priority jobs from occupying every slot once they are running; it only decides which waiting job starts next. If slow jobs already hold all slots, urgent jobs still wait for one to finish. That is why priority complements per-type limits rather than replacing them.
Verify: with all slots busy on long jobs, measure urgent-job latency; if it equals the long jobs' remaining duration, you need a limit or a separate queue, not just priority.
Verification¶
Per-type limits are correct when:
- No job waits for a limit while holding a worker slot.
- Latency-sensitive job types keep their latency while heavy types are backlogged.
- Fleet-wide limits are enforced with leases that expire when workers crash.
- Deferral tries are budgeted separately from real failures.
Diagnostic Hook: export in-progress jobs per type per worker and queue wait per type. A heavy type holding most slots while a light type's wait rises is the starvation in this guide; a high deferral rate for one type means its limit is lower than its arrival rate, so its backlog will keep growing until the limit or the capacity changes.
Pitfalls & edge cases¶
- Semaphores inside jobs. Waiting jobs occupy slots and starve everything else.
- Deferrals counted as failures. Jobs exhaust
max_trieswhile merely waiting their turn. - Leases shorter than jobs. A running job loses its slot and the limit is exceeded.
- Separate queues with shared workers. If one worker polls both queues with one
max_jobs, the isolation is lost.
Frequently Asked Questions¶
How do I limit how many jobs of one type run at once?
Either run that job type on its own queue with its own workers and max_jobs, or have the job defer itself with arq.Retry when a per-type semaphore or Redis counter is at its limit. Do not wait on a semaphore inside the job.
Why did adding a semaphore to my job make other jobs slower?
Jobs waiting on the semaphore still occupy worker slots, so fewer slots are left for other job types. In testing, email latency rose from 1.5 s to 3.7 s when reports were limited this way.
How do I enforce a concurrency limit across many workers?
Keep a shared count in Redis — for example a sorted set of leases scored by expiry — acquire atomically with a Lua script, defer the job when the limit is reached, and release in a finally block. Expiring leases free slots held by crashed workers.
Is job priority enough to protect urgent jobs?
Not on its own. Priority picks the next job to start, but if slow jobs already hold every slot, urgent jobs still wait for one to finish. Combine priority with per-type limits or separate queues.
Related¶
- Background Jobs & Task Queues — up to the topic overview.
- Isolating tenants with bulkheads — the same idea applied per customer.
- Concurrent Execution & Worker Patterns — the section overview.