Skip to content

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

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.

Email latency while 30 slow reports are queued 4 horizontal bars comparing semaphore inside the job with the others. Email latency while 30 slow reports are queued semaphore inside the job 3,706 ms shared queue, no limit 1,545 ms Retry(defer) when busy 40 ms separate queues 36 ms arq 0.28, max_jobs=10 total; 30 reports of 500 ms and 20 emails of 10 ms. Limiting work without freeing slots starves everything else; free the slot or split the queue.

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.

Ways to cap one job type's concurrency A grid of 4 rows by 4 columns. Ways to cap one job type's concurrency approach frees the slot limit scope cost semaphore inside job no per worker starves other types Retry(defer) when busy yes per worker extra tries and queue traffic separate queues n/a: own slots per worker pool more deployments Redis counter + defer yes whole fleet a Redis round trip per job Separate queues are the default; deferral fits when splitting deployments is not worth it.

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.

How should this job type be isolated? A decision on What must the limit protect with 3 outcomes. How should this job type be isolated? What must the limit protect? other job types' latency separate queues own workers and max_jobs a resource on one worker defer when busy Retry(defer=...) an external quota Redis lease set fleet-wide count Never limit by waiting inside a slot; free the slot or give the job type its own.

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