Running arq Workers with Redis¶
arq is the background-job library that fits asyncio services most naturally: jobs are async def functions, a worker is a single event loop running many jobs concurrently, and Redis is the only dependency. That concurrency is also its main tuning knob. Measured with arq 0.28 against Redis 7.4, a burst worker processing 2,000 jobs that each awaited 10 ms of I/O managed 474 jobs/s with max_jobs=10, 2,251 jobs/s with max_jobs=50, and 2,367 jobs/s with max_jobs=200 — after the first jump, the queue and Redis round trips, not concurrency, set the limit. This guide sets up a production worker: settings, shared clients through startup hooks, the right max_jobs, scheduled jobs, health checks, and a clean stop.
Prerequisites¶
- Python 3.11+,
pip install arqand a Redis server. - How arq compares, from choosing between Celery, arq and taskiq.
- Client lifetimes, from reusing aiohttp ClientSession across requests.
1. Define jobs and WorkerSettings¶
A worker is configured by a class whose attributes arq reads; run it with the arq command:
# worker.py (pip install arq)
import httpx
from arq.connections import RedisSettings
from arq.cron import cron
async def startup(ctx: dict) -> None:
ctx["http"] = httpx.AsyncClient(timeout=10) # one client for every job
ctx["db"] = await create_db_pool()
async def shutdown(ctx: dict) -> None:
await ctx["http"].aclose()
await ctx["db"].close()
async def send_webhook(ctx: dict, url: str, payload: dict) -> int:
r = await ctx["http"].post(url, json=payload)
r.raise_for_status()
return r.status_code
class WorkerSettings:
functions = [send_webhook]
on_startup = startup
on_shutdown = shutdown
redis_settings = RedisSettings(host="redis", port=6379)
max_jobs = 50
job_timeout = 60
keep_result = 3600
arq worker.WorkerSettings
ctx is a dict shared by every job on this worker. Creating clients and pools in on_startup and closing them in on_shutdown gives each worker one connection pool instead of one per job — the same rule as for web servers, with the same leak if ignored.
Verify: start the worker, enqueue a job from a REPL, and see it complete in the worker log.
2. Enqueue jobs from the web service¶
Producers use an ArqRedis pool, typically created in the web app's lifespan:
from contextlib import asynccontextmanager
from arq import create_pool
from fastapi import FastAPI
@asynccontextmanager
async def lifespan(app: FastAPI):
app.state.arq = await create_pool(RedisSettings(host="redis"))
yield
await app.state.arq.aclose()
app = FastAPI(lifespan=lifespan)
@app.post("/orders/{order_id}/notify")
async def notify(order_id: str):
job = await app.state.arq.enqueue_job(
"send_webhook", WEBHOOK_URL, {"order": order_id},
_job_id=f"notify:{order_id}", # duplicate enqueues are ignored
)
return {"queued": job is not None}
_job_id makes the enqueue itself idempotent: verified, a second enqueue_job with an existing job id returned None instead of queueing a duplicate. That handles a client retrying the HTTP request; the job body still needs its own idempotency for redeliveries, as in making background jobs idempotent.
Verify: call the endpoint twice for one order; one job is queued.
3. Size max_jobs from the jobs' I/O profile¶
max_jobs is how many jobs one worker runs concurrently. For I/O-bound jobs the right value is large, because each job spends most of its time awaiting:
class WorkerSettings:
functions = [send_webhook]
max_jobs = 50 # measured knee for 10 ms I/O jobs
The measurement for 2,000 jobs of 10 ms: 4.22 s at max_jobs=10, 0.89 s at 50, 0.84 s at 200. Beyond the knee, extra concurrency mostly adds pressure on whatever the jobs call — the webhook targets, the database pool — without adding throughput. If jobs share a pool of 20 database connections, max_jobs above 20 just queues jobs on the pool; size them together, as discussed in sizing async connection pools for throughput.
CPU-heavy jobs are different: one CPU-bound job blocks every other job on the worker's loop. Send that work to a process pool from inside the job, or run a separate worker deployment with max_jobs=1 per process.
Verify: benchmark your real job mix at three max_jobs values and pick the knee.
4. Schedule recurring jobs with cron¶
arq runs scheduled jobs from the worker itself, with second-level precision and a uniqueness lock so that several workers do not all run the same scheduled job:
from arq.cron import cron
async def prune_processed_jobs(ctx: dict) -> int:
return await ctx["db"].fetchval("select prune_processed(14)")
class WorkerSettings:
functions = [send_webhook]
cron_jobs = [
cron(prune_processed_jobs, hour={3}, minute={17}, run_at_startup=False),
]
By default each cron job is unique across workers sharing the queue, so scaling to five workers does not run the prune five times. Pick odd minutes (minute={17}) rather than :00, so your jobs do not coincide with everyone else's hourly load. In-process scheduling without a job system is covered in scheduling cron jobs inside an asyncio service.
Verify: with three workers running, the scheduled job runs once per schedule, visible in the logs of exactly one worker.
5. Health-check and stop the worker cleanly¶
arq periodically writes a health-check record to Redis, and arq --check reads it — the basis for a container liveness probe:
arq --check worker.WorkerSettings # exits 0 if a recent health check exists
Set health_check_interval shorter than the probe period. Shutdown needs a deliberate choice: by default, on SIGTERM arq cancels running jobs immediately — they are retried later by another worker, counting against max_tries — rather than letting them finish. Reading arq 0.28's signal handler confirms it: every unfinished job task is cancelled. To let in-flight jobs complete, set job_completion_wait (seconds to wait before cancelling), and give the container a termination grace period longer than that:
class WorkerSettings:
functions = [send_webhook]
job_timeout = 60
job_completion_wait = 45 # on SIGTERM: stop pulling, let running jobs finish for up to 45 s
# Kubernetes container spec (excerpt)
terminationGracePeriodSeconds: 60 # > job_completion_wait of 45
livenessProbe:
exec:
command: ["arq", "--check", "worker.WorkerSettings"]
periodSeconds: 60
The general shutdown sequence and how it interacts with Kubernetes is in shutting down asyncio pods in Kubernetes.
Verify: send SIGTERM during a batch; with job_completion_wait set, running jobs complete, no new ones start, and the process exits within the grace period; without it, they are cancelled and show up again as retries.
Verification¶
The worker is production-ready when:
- Shared clients live in
ctx, created inon_startupand closed inon_shutdown. max_jobsis chosen from a measurement, consistent with downstream pool sizes.- Enqueues use
_job_idwhere producers can repeat, and jobs are idempotent. - Health checks,
job_completion_waitand termination grace are set together.
Diagnostic Hook: export queued jobs (the length of arq's queue in Redis), jobs in progress per worker, job duration by function, and failures by function. A queue that grows while in-progress sits at max_jobs on every worker means you need more workers; one that grows while in-progress is low means workers are not pulling — check their health-check records and logs.
Pitfalls & edge cases¶
- CPU-bound work in a job. It blocks every other job on that worker's loop.
- Creating clients per job. One connection pool per job is a descriptor leak under load.
- Relying on the default SIGTERM behaviour. Running jobs are cancelled and retried on every deploy unless
job_completion_waitis set. - Results kept forever.
keep_resultcontrols how long results stay in Redis; large results with long retention fill memory.
Frequently Asked Questions¶
How do I share an HTTP client or database pool across arq jobs?
Create it in the worker's on_startup hook and store it in ctx, then close it in on_shutdown. Every job receives the same ctx dict.
What should arq max_jobs be?
For I/O-bound jobs, high enough to keep the worker busy while jobs await — in testing, throughput rose from 474 to 2,251 jobs per second between 10 and 50 and barely moved after. Keep it consistent with downstream pool sizes.
How do I prevent duplicate arq jobs?
Pass _job_id when enqueueing; a second enqueue with an existing job id returns None instead of queueing another job. The job itself should still be idempotent for redeliveries.
Do arq cron jobs run on every worker?
No. By default each cron job is unique across workers sharing a queue, so it runs once per schedule even with several workers.
Related¶
- Background Jobs & Task Queues — up to the topic overview.
- Limiting concurrency per job type — when one job type must not use all of max_jobs.
- Concurrent Execution & Worker Patterns — the section overview.