Guide 05
Example: Celery + RabbitMQ + FastAPI
Audience: Teams standardizing on Celery for durable jobs. Broker: RabbitMQ (AMQP). Redis only as optional result backend.
TL;DR layout
text
app/
core/settings.py
api/routes/jobs.py
workers/celery_app.py
workers/tasks/email.py
workers/tasks/exports.py
Dockerfile # same image, different CMD
docker-compose.yml # api, worker, rabbitmq, db, flower?Contents
- Settings
- Celery app
- Task definition
- FastAPI enqueue
- Queues and routing
- Retries and time limits
- Run commands
- Ops notes
---
1. Settings
python
# PSEUDOCODE : pydantic-settings
class Settings(BaseSettings):
database_url: str
celery_broker_url: str = "amqps://user:pass@rabbitmq:5671//"
celery_result_backend: str | None = None # optional Redis/DB; status still in app DB
environment: str = "local"
model_config = SettingsConfigDict(env_file=".env", extra="ignore")2. Celery app
python
# PSEUDOCODE : app/workers/celery_app.py
from celery import Celery
from app.core.settings import settings
celery_app = Celery("app", broker=settings.celery_broker_url)
celery_app.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
task_acks_late=True,
worker_prefetch_multiplier=1, # fair for long jobs
task_default_queue="jobs.default",
task_routes={
"app.workers.tasks.exports.*": {"queue": "jobs.heavy"},
"app.workers.tasks.email.*": {"queue": "jobs.io"},
},
broker_connection_retry_on_startup=True,
)
celery_app.autodiscover_tasks(["app.workers.tasks"])Never enable pickle.
3. Task definition
python
# PSEUDOCODE : app/workers/tasks/email.py
from app.workers.celery_app import celery_app
from app.db import session_scope
from app.services.mail import send_email
@celery_app.task(
bind=True,
name="app.workers.tasks.email.send_receipt",
max_retries=5,
autoretry_for=(TransientMailError,),
retry_backoff=True,
retry_backoff_max=600,
retry_jitter=True,
soft_time_limit=30,
time_limit=45,
)
def send_receipt(self, job_id: str):
with session_scope() as db:
job = db.get_job(job_id)
if job is None or job.status == "succeeded":
return
db.mark_running(job_id)
try:
send_email(job.payload_ref) # sync client OK in Celery prefork
db.mark_succeeded(job_id)
except PermanentMailError as e:
db.mark_failed(job_id, str(e))
raise
except TransientMailError as e:
db.bump_attempt(job_id, str(e))
raise self.retry(exc=e)4. FastAPI enqueue
python
# PSEUDOCODE
@router.post("/receipts", status_code=202)
async def enqueue_receipt(body: ReceiptIn, db: Db = Depends()):
job = await db.create_job(type="send_receipt", payload_ref=..., idempotency_key=body.key)
send_receipt.delay(str(job.id)) # or apply_async(queue="jobs.io")
return {"job_id": job.id, "status": "pending"}Prefer job_id only on the wire; load details from DB in the worker.
5. Queues and routing
| Queue | Work |
|---|---|
jobs.default | General |
jobs.high | User-visible latency |
jobs.heavy | CPU / large files |
jobs.io | Email / HTTP |
jobs.dlq | Dead letters (via DLX) |
Separate worker deployments can subscribe to different queues.
6. Retries and time limits
acks_late=Trueso crash redelivers- Soft + hard time limits
- Retry only transient errors
- After max retries → mark failed + DLQ
See 09 Error taxonomy.
7. Run commands
bash
# API
uvicorn app.main:app --host 0.0.0.0 --port 8000
# Worker
celery -A app.workers.celery_app.celery_app worker -Q jobs.default,jobs.io -c 4
# Heavy pool
celery -A app.workers.celery_app.celery_app worker -Q jobs.heavy -c 1
# Optional Beat (single leader)
celery -A app.workers.celery_app.celery_app beat
# Optional Flower (private + auth)
celery -A app.workers.celery_app.celery_app flower8. Ops notes
- Memory: recycle workers (
worker_max_tasks_per_child) - Graceful shutdown: stop consuming, finish in-flight within termination grace
- Metrics: task runtime, retries, queue depth via RabbitMQ
- Flower is not the product job status UI
Related: 06 Taskiq if you want async-native workers.
Also: rate limits + entity locks end-to-end : 17 Celery RL + lock pipeline