Background Jobs: Celery Architecture, Retries & Task Canvas
Offloading heavy, long-running computations (sending emails, video encoding, PDF generation, machine learning inference) out of HTTP request-response cycles requires a distributed background task queue. In Python, the industry standard is Celery.
This chapter details Celery architecture (Brokers, Result Backends, Worker processes), Exponential Backoff Retries with Jitter, Task Idempotency, Dead Letter Queues (DLQ), and Celery Task Canvas (chain, group, chord).
1. Celery Distributed Architecture
Celery decouples web applications from background task execution using a Message Broker:
Celery Architecture Component Hierarchy:
[ Producer (FastAPI / Django) ] ββ> Calls: send_email.delay(user_id)
|
v (Serializes task message into Queue)
[ Message Broker (Redis / RabbitMQ) ]
|
v (Pulls task message from queue)
[ Celery Worker Pool (Prefork / Gevent) ]
|
βββ Executes Task Function
v (Optional: Stores task state/return value)
[ Result Backend (Redis / PostgreSQL) ]2. Exponential Backoff Retries with Jitter
Network requests inside background tasks frequently fail due to temporary microservice outages or rate limits.
CRITICAL RETRY INVARIANT: Always use Exponential Backoff with Jitter to prevent Thundering Herd retries!
# app/tasks.py
from celery import SharedTask
import logging
import random
import requests
logger = logging.getLogger(__name__)
@shared_task(
bind=True,
max_retries=5,
default_retry_delay=2,
autoretry_for=(requests.RequestException,),
retry_backoff=True, # Enables Exponential Backoff (2s -> 4s -> 8s -> 16s...)
retry_backoff_max=600, # Cap backoff at 10 minutes max
retry_jitter=True # Adds random jitter to prevent thundering herd retries!
)
def send_webhook(self, url: str, payload: dict):
logger.info("Executing webhook task attempt %s", self.request.retries)
response = requests.post(url, json=payload, timeout=10)
response.raise_for_status()
return response.json()3. Task Idempotency Invariants
Celery worker processes can crash or experience network partitions right after completing a task but before sending task acknowledgment (ACK) to the Broker.
If a task is re-queued and executed a second time, it must be Idempotent (producing identical business results regardless of execution count):
# Idempotent Task Design using Redis Mutex Key
from app.celery_app import cel_app
import redis
redis_client = redis.Redis()
@cel_app.task(bind=True)
def process_payment(self, payment_id: str, amount: float):
lock_key = f"payment_processed:{payment_id}"
# Atomic SET NX (Set if Not Exists)
if not redis_client.set(lock_key, "1", nx=True, ex=86400):
logger.warning("Payment %s already processed! Skipping idempotent task execution.", payment_id)
return {"status": "already_processed"}
# Execute actual payment transaction safely...
charge_credit_card(payment_id, amount)
return {"status": "success"}4. Complex Workflow Workflows: Celery Canvas (chain, group, chord)
Celery Canvas provides functional primitives to compose complex task execution DAGs:
-
chain(Sequential Execution): Executes tasks sequentially, passing the result of Task 1 as the first argument to Task 2:from celery import chain # Workflow: fetch_data() -> transform() -> save() workflow = chain(fetch_data.s(url), transform.s(), save.s()) workflow.apply_async() -
group(Parallel Execution): Executes multiple tasks concurrently in parallel across workers:from celery import group # Runs 3 image resize tasks in parallel! job = group(resize_image.s(img_id) for img_id in image_ids) job.apply_async() -
chord(Parallel Execution + Callback Synchronization): Executes agroupof tasks in parallel, and passes all aggregated results to a final callback task once all group tasks finish:from celery import chord # Executes 10 sub-tasks in parallel, then passes list of 10 results to generate_report! workflow = chord( (fetch_part.s(i) for i in range(10)), generate_report.s() # Callback function ) workflow.apply_async()