Celery Patterns
Setup
# celery_app.py
from celery import Celery
from kombu import Queue
celery_app = Celery(
"myapp",
broker="redis://redis:6379/0",
backend="redis://redis:6379/1",
include=["app.tasks.email", "app.tasks.reports"],
)
celery_app.conf.update(
task_serializer="json",
result_serializer="json",
accept_content=["json"],
timezone="UTC",
enable_utc=True,
# Queues
task_queues=[
Queue("high", routing_key="high"),
Queue("default", routing_key="default"),
Queue("low", routing_key="low"),
],
task_default_queue="default",
# Reliability
task_acks_late=True, # ack after completion, not on receipt
task_reject_on_worker_lost=True,
worker_prefetch_multiplier=1, # prevent worker starvation
# Result expiry
result_expires=3600,
)
Task Definition
from celery import shared_task
from celery.utils.log import get_task_logger
from celery.exceptions import MaxRetriesExceededError
logger = get_task_logger(__name__)
@shared_task(
bind=True,
max_retries=3,
default_retry_delay=60,
queue="default",
time_limit=300, # hard kill after 5 min
soft_time_limit=240, # raises SoftTimeLimitExceeded after 4 min
autoretry_for=(TemporaryError,),
retry_backoff=True,
retry_backoff_max=600,
retry_jitter=True,
)
def send_email_task(self, to: str, subject: str, body: str) -> dict:
try:
result = email_service.send(to=to, subject=subject, body=body)
return {"status": "sent", "message_id": result.message_id}
except PermanentError as exc:
logger.error("Permanent failure sending to %s: %s", to, exc)
raise # don't retry permanent errors
except TemporaryError as exc:
logger.warning("Temporary failure, will retry: %s", exc)
raise self.retry(exc=exc)
Calling Tasks
# Fire and forget
send_email_task.delay(to="user@example.com", subject="Welcome", body="...")
# With options
send_email_task.apply_async(
kwargs={"to": "user@example.com", "subject": "...", "body": "..."},
countdown=60, # delay 60 seconds
expires=3600, # discard if not started within 1h
queue="high", # override queue
)
# Chain (pipeline)
from celery import chain
result = chain(
fetch_data.s(source_id),
process_data.s(),
save_results.s(),
).delay()
# Group (parallel)
from celery import group
result = group(
send_email_task.s(email, "Welcome", body) for email in email_list
).delay()
Periodic Tasks (Celery Beat)
from celery.schedules import crontab
celery_app.conf.beat_schedule = {
"daily-report": {
"task": "app.tasks.reports.generate_daily_report",
"schedule": crontab(hour=8, minute=0),
"kwargs": {"format": "pdf"},
},
"cleanup-expired-sessions": {
"task": "app.tasks.cleanup.remove_expired_sessions",
"schedule": 300.0, # every 5 minutes
},
}
Monitoring
# Flower (web UI)
celery -A celery_app flower --port=5555
# CLI inspection
celery -A celery_app inspect active # running tasks
celery -A celery_app inspect reserved # queued tasks
celery -A celery_app inspect stats # worker stats
celery -A celery_app events # real-time event stream