How to Set Up and Configure Celery for Asynchronous Task Processing in PostHog
PostHog configures Celery using Redis as both broker and result backend, initializes the application in posthog/celery.py with Prometheus metrics and auto-discovered tasks, and routes work to named queues defined in bin/celery-queues.env.
PostHog relies on Celery to handle background operations like cohort recalculations, email dispatching, and analytics queries without blocking the main web process. To extend or debug this pipeline, you must understand how the Django application wires together the message broker, task queues, and worker processes. This guide provides exact steps to set up and configure Celery for asynchronous task processing based on the production implementation in the PostHog/posthog repository.
Configure Redis as the Message Broker and Result Backend
All Celery configuration resides in posthog/settings/celery.py. By default, PostHog uses Redis for both task queuing and result storage, reading from the shared REDIS_URL environment variable:
# posthog/settings/celery.py
CELERY_BROKER_URL = REDIS_URL # → redis://<host>:6379/0
CELERY_RESULT_BACKEND = REDIS_URL # same Redis instance
If your deployment requires a dedicated Redis instance for Celery traffic, override REDIS_URL in the worker environment (e.g., REDIS_URL=redis://celery-redis:6379/1). For test environments, the configuration automatically switches to eager mode when the TEST flag is true, executing tasks synchronously:
if TEST:
import celery
celery.current_app.conf.CELERY_ALWAYS_EAGER = True
celery.current_app.conf.CELERY_EAGER_PROPAGATES_EXCEPTIONS = True
Initialize the Celery Application
The posthog/celery.py file creates a singleton Celery instance named "posthog" and hooks it into Django settings:
# posthog/celery.py
app = Celery("posthog")
app.config_from_object("django.conf:settings", namespace="CELERY")
app.autodiscover_tasks() # loads tasks.py from every Django app
This initialization includes several production-grade features:
- Prometheus metrics – Signal handlers increment counters for
pre_run,success,failure, andretryevents (lines 30–62) - Worker process initialization – Starts an HTTP endpoint on
CELERY_METRICS_PORT(default 8001) to expose metrics and calls_initialize_worker_metrics()for long-running workers (lines 36–45, 138–146) - ClickHouse tagging – Automatically tags all queries executed inside Celery tasks as
celeryworkload to isolate them from online traffic (lines 148–159) - Periodic task registration – After the app finalizes,
setup_periodic_tasksfromposthog/tasks/scheduled.pyregisters crontab and interval jobs (lines 194–204)
Define Named Task Queues
PostHog workers subscribe to specific named queues rather than the default celery queue alone. The authoritative list lives in bin/celery-queues.env:
# bin/celery-queues.env
CELERY_WORKER_QUEUES=celery,stats,email,analytics_queries,analytics_limited,long_running,exports,subscription_delivery,usage_reports,integrations,feature_flags,feature_flags_long_running
When you introduce a new queue, you must synchronize two locations:
- Add the queue name to
bin/celery-queues.env - Add a corresponding entry in
posthog/tasks/utils.py(which defines theCeleryQueueenum used by worker startup logic)
Writing Asynchronous Tasks
Task modules import the shared_task decorator from Celery. The following pattern from products/visual_review/backend/tasks/tasks.py demonstrates best practices:
# products/visual_review/backend/tasks/tasks.py
from celery import shared_task
@shared_task(ignore_result=True) # avoids unnecessary result storage
def process_visual_review(review_id: int) -> None:
# implementation logic here
pass
For long-running operations, specify a time_limit and raise exceptions to trigger the built-in retry mechanism:
from celery import shared_task
@shared_task(ignore_result=True, time_limit=300)
def recalculate_team_cohorts(team_id: int) -> None:
"""Heavy calculation that runs for up to 5 minutes."""
from posthog.models import Team
team = Team.objects.get(id=team_id)
# perform calculation...
PostHog recommends ignore_result=True (or the global CELERY_IGNORE_RESULT setting) to prevent Redis memory bloat on high-volume tasks.
Run Celery Workers Locally
Start a development worker by first exporting the queue definition and then executing the Celery command:
# Load queue list into the environment
export $(cat bin/celery-queues.env | xargs)
# Start a worker processing specific queues
celery -A posthog worker -Q celery,long_running --loglevel=INFO
To process all configured queues at once:
celery -A posthog worker -Q $(echo $CELERY_WORKER_QUEUES) --loglevel=INFO --concurrency=4
Deploy Celery Workers in Production
In Docker Compose or Kubernetes deployments, mount the environment file and pass the queue variable to the worker command:
# docker-compose.yml snippet
services:
celery-worker:
image: posthog/posthog:latest
command: celery -A posthog worker -Q $CELERY_WORKER_QUEUES --loglevel=INFO
env_file:
- bin/celery-queues.env
- .env
depends_on:
- redis
The worker container shares the same Redis instance as the web process, though you may allocate a dedicated Redis cluster for higher isolation between web and background workloads.
Monitor Celery Performance and Health
PostHog exposes task metrics via a Prometheus HTTP endpoint running on the port defined by CELERY_METRICS_PORT (default 8001). Key metrics include:
posthog_celery_task_total– Counter for task lifecycle events (success, failure, retry)posthog_celery_task_duration_seconds– Histogram of execution times
Add this scrape configuration to your Prometheus setup:
- job_name: posthog_celery
static_configs:
- targets: ['celery-worker:8001']
Additionally, the utility function posthog.utils.get_celery_heartbeat() returns either "offline" or a timestamp, enabling health checks in load balancers or container orchestrators.
Summary
- Broker configuration –
posthog/settings/celery.pysetsCELERY_BROKER_URLandCELERY_RESULT_BACKENDtoREDIS_URL, with eager mode automatically enabled during tests. - App initialization –
posthog/celery.pycreates the Celery singleton, auto-discovers tasks, installs Prometheus metrics, and registers periodic tasks viasetup_periodic_tasks. - Queue management – Worker queue names are defined in
bin/celery-queues.envand must stay synchronized withposthog/tasks/utils.py. - Task implementation – Use
@shared_task(ignore_result=True)to define background jobs; metrics and ClickHouse tagging are handled automatically. - Deployment – Run workers with
celery -A posthog worker -Q $CELERY_WORKER_QUEUESand expose port 8001 for Prometheus scraping.
Frequently Asked Questions
What Redis URL does PostHog Celery use?
Both the broker and result backend read from the REDIS_URL environment variable, as configured in posthog/settings/celery.py. If unset, the application defaults to redis://localhost:6379/0. You can override this per environment by setting REDIS_URL to a dedicated Redis instance for Celery isolation.
How do I add a new task queue to PostHog?
First append the queue name to bin/celery-queues.env in the CELERY_WORKER_QUEUES comma-separated list. Then add a corresponding entry to the CeleryQueue enum in posthog/tasks/utils.py. Finally, restart your workers to begin consuming the new queue.
How do I run Celery tasks synchronously for testing?
When the TEST environment variable is set to true, posthog/settings/celery.py automatically configures CELERY_ALWAYS_EAGER = True and CELERY_EAGER_PROPAGATES_EXCEPTIONS = True. This forces all .delay() and .apply_async() calls to execute inline, making debugging and unit testing straightforward.
Where are periodic Celery tasks registered in PostHog?
The setup_periodic_tasks function in posthog/tasks/scheduled.py registers all crontab and interval schedules. This function is called via the worker_ready signal defined in posthog/celery.py (lines 194–204), ensuring periodic tasks are configured every time a worker process starts.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →