Testing Celery Tasks with pytest

Celery ships its own pytest plugin and several ways to shortcut the broker, and choosing between them is most of what makes a Celery test suite trustworthy. This guide builds a layered suite with pytest as part of Testing Background Jobs in Backend Frameworks & Worker Scaling: direct task calls for logic, captured enqueues for contracts, and an embedded worker against real Redis for chords, retries, and serialization.

Problem Statement

A Django project's Celery tests all run with CELERY_TASK_ALWAYS_EAGER = True in the test settings. The suite is green, yet production has repeatedly hit bugs the tests could not have caught: a chord that never fired, a task that failed to serialize a Decimal argument, and retry settings that retried a permanent validation error. You want tests that call task logic cheaply, verify that and when tasks are enqueued without running them, exercise the real retry machinery, and run workflows through a real broker and worker — all within a CI budget of a few minutes.

Prerequisites

  • Celery 5.3+ and pytest 7+; the plugin is enabled with pytest -p celery.contrib.pytest or pytest_plugins = ("celery.contrib.pytest",) in conftest.py.
  • Docker available in CI for the integration layer (testcontainers or a service container).
  • Tasks written so their logic can be called as plain functions (keep task bodies thin).
  • freezegun or time-machine for tests involving time.

Step 1 — Turn Eager Mode Off

Eager mode (task_always_eager) runs .delay() synchronously in the caller. It bypasses serialization, routing, the result backend, acknowledgement, and real retries. Remove it from test settings and replace its two legitimate uses — "run the logic" and "check it was enqueued" — with explicit tools.

# settings/test.py
CELERY_TASK_ALWAYS_EAGER = False          # was True: hid serialization and chord bugs
CELERY_TASK_EAGER_PROPAGATES = True       # only relevant if someone re-enables eager locally
CELERY_BROKER_URL = "memory://"           # default for tests that don't start a worker

Tests that relied on eager mode now fail loudly, which is the point: each one is either testing logic (Step 2) or testing an enqueue (Step 3), and should say which.

Replace eager mode with three explicit tools Calling task.run executes the task body as a plain function for logic tests. Capturing apply_async records what was enqueued without running it, for contract tests. The celery_worker fixture runs an embedded worker against a real broker for serialization, retries, and chords. One tool per question task.run(...) does the logic work? milliseconds capture apply_async was it enqueued, when, with which args? celery_worker does it survive the queue? seconds, real Redis Eager mode tried to answer all three and answered the last one wrongly.

Step 2 — Test Logic by Calling the Task Body

task.run(*args) calls the undecorated function. For bound tasks (bind=True), run supplies self automatically; to control self.request (retries, id), use task.apply() with explicit options or push a request context.

# tasks.py
@app.task(bind=True, autoretry_for=(ConnectionError,), retry_backoff=True, max_retries=5)
def sync_customer(self, customer_id: int) -> str:
    customer = Customer.objects.get(pk=customer_id)
    if customer.synced_version == customer.version:
        return "noop"                               # idempotent: nothing changed
    crm.upsert(customer)
    Customer.objects.filter(pk=customer_id).update(synced_version=customer.version)
    return "synced"

# test_tasks.py
@pytest.mark.django_db
def test_sync_is_idempotent(crm_fake):
    c = Customer.objects.create(name="Ada", version=3, synced_version=2)
    assert sync_customer.run(c.id) == "synced"
    assert sync_customer.run(c.id) == "noop"        # second delivery changes nothing
    assert crm_fake.upserts == 1

@pytest.mark.django_db
def test_sync_sees_retry_count():
    c = Customer.objects.create(name="Ada", version=1, synced_version=0)
    result = sync_customer.apply(args=[c.id], retries=4)    # simulate 5th attempt
    assert result.successful()

apply() runs the task locally including its request context and retry handling, but still without serialization or a broker. Use it for logic that depends on self.request.retries.

Step 3 — Capture Enqueues Instead of Running Them

To assert that code enqueues a task — and only after the transaction commits — patch apply_async with a recorder. With Django, transaction.on_commit callbacks do not run inside TestCase transactions by default; use django_capture_on_commit_callbacks (pytest-django) or TransactionTestCase semantics.

# conftest.py
@pytest.fixture
def enqueued(monkeypatch):
    calls = []
    def fake_apply_async(self, args=None, kwargs=None, **options):
        calls.append({"task": self.name, "args": args or (), "kwargs": kwargs or {}, **options})
        return AsyncResult("fake-id")
    monkeypatch.setattr(celery.app.task.Task, "apply_async", fake_apply_async)
    return calls

# test_signup.py
@pytest.mark.django_db
def test_welcome_email_enqueued_after_commit(enqueued, django_capture_on_commit_callbacks):
    with django_capture_on_commit_callbacks(execute=False) as callbacks:
        signup(email="ada@example.com")
    assert enqueued == []                               # nothing before commit
    for cb in callbacks:
        cb()                                            # simulate the commit
    assert [c["task"] for c in enqueued] == ["accounts.tasks.send_welcome_email"]
    assert enqueued[0]["queue"] == "email"              # routing contract

@pytest.mark.django_db
def test_no_email_on_rollback(enqueued, django_capture_on_commit_callbacks):
    with django_capture_on_commit_callbacks(execute=True):
        with pytest.raises(ValidationError):
            signup(email="not-an-email")
    assert enqueued == []

This is the test that proves the enqueue follows the transaction — the property the transactional outbox pattern exists to guarantee.

The enqueue follows the transaction The test captures on-commit callbacks. While the signup transaction is open, the recorder is empty. Running the commit callbacks records send_welcome_email routed to the email queue. In a second test the signup raises and rolls back, and the recorder stays empty. What the recorder sees inside transaction recorder: [] on commit send_welcome_email, queue = email rollback: stays []

Step 4 — Test the Retry Configuration

autoretry_for, retry_backoff, and max_retries are configuration, and configuration deserves tests. Assert which exceptions trigger a retry and that the delay sequence matches what you intend.

from celery.exceptions import Retry

def test_connection_error_retries(monkeypatch):
    monkeypatch.setattr(crm, "upsert", Mock(side_effect=ConnectionError))
    with pytest.raises(Retry):
        sync_customer.apply(args=[customer.id], throw=True).get()

def test_validation_error_does_not_retry(monkeypatch):
    monkeypatch.setattr(crm, "upsert", Mock(side_effect=CrmValidationError("bad phone")))
    result = sync_customer.apply(args=[customer.id])
    assert result.failed() and isinstance(result.result, CrmValidationError)

def test_backoff_sequence_is_capped():
    from celery.utils.time import get_exponential_backoff_interval
    delays = [get_exponential_backoff_interval(factor=1, retries=n, maximum=600, full_jitter=False)
              for n in range(12)]
    assert delays[:4] == [1, 2, 4, 8] and max(delays) == 600

Checking the delay helper with jitter disabled verifies the curve; production keeps jitter on. The reasoning behind these settings is in exponential backoff with jitter in Celery, and a framework-neutral approach is in testing retry and backoff logic deterministically.

Step 5 — Run Workflows Through a Real Worker

The plugin's celery_app and celery_worker fixtures start an embedded worker in a thread. Point them at real Redis to test chords, serialization, and routing.

# conftest.py
from testcontainers.redis import RedisContainer

@pytest.fixture(scope="session")
def redis_container():
    with RedisContainer("redis:7.2") as r:
        yield r

@pytest.fixture(scope="session")
def celery_config(redis_container):
    url = f"redis://{redis_container.get_container_host_ip()}:{redis_container.get_exposed_port(6379)}"
    return {"broker_url": f"{url}/0", "result_backend": f"{url}/1",
            "task_serializer": "json", "accept_content": ["json"], "task_acks_late": True}

@pytest.fixture(scope="session")
def celery_worker_parameters():
    return {"queues": ("default", "email", "reports"), "perform_ping_check": False}

# test_workflows.py
@pytest.mark.integration
def test_statement_chord_fires_body(celery_app, celery_worker):
    res = statement_workflow("c-17", "2026-08").apply_async()
    assert res.get(timeout=30) is None
    assert storage.exists("pdf/c-17.pdf")

@pytest.mark.integration
def test_decimal_argument_fails_to_serialize(celery_worker):
    with pytest.raises(EncodeError):
        charge.delay(Decimal("9.99"))                   # json serializer rejects Decimal

The serialization test documents a real constraint: with the JSON serializer, pass amounts as strings or integer cents. Tasks under test must be imported so the embedded worker registers them; add celery_includes if they live in modules the test does not import.

Real broker, embedded worker The pytest process runs the test and an embedded worker thread started by the celery_worker fixture. Both talk to a Redis container: database 0 is the broker and database 1 is the result backend. Tasks are serialized as JSON and go through the broker, so chords and routing behave as in production. Integration layer topology pytest process test code worker thread celery_worker Redis container db 0: broker db 1: result backend

Step 6 — Keep the Integration Layer Fast and Isolated

Session-scoped containers and workers are fast but share state between tests. Isolate with a flush between tests and unique identifiers per test, and mark the layer so it runs as a separate CI job.

@pytest.fixture(autouse=True)
def clean_redis(request, redis_container):
    if "integration" in request.keywords:
        redis_container.get_client().flushall()
    yield

# pytest.ini
# [pytest]
# markers = integration: needs Docker and a real broker
# addopts = -m "not integration"          # default local run: unit + contract only

CI then runs pytest (unit and contract, parallel with -n auto) and pytest -m integration as two jobs. Retry delays for integration tests can be shortened through settings (retry_backoff_max, custom countdown multipliers) so real retries complete in seconds.

Step 7 — Guard Routing and Task Names with a Configuration Test

Two production failures come from configuration drift rather than code: a task renamed or moved to another module (so messages already queued under the old name fail with NotRegistered), and a routing rule that silently stops matching after a refactor (so a heavy task lands on the latency-sensitive default queue). Both are cheap to pin with a test that inspects the app instead of running anything.

EXPECTED_ROUTES = {
    "accounts.tasks.send_welcome_email": "email",
    "billing.tasks.charge_invoice": "billing",
    "reports.tasks.build_monthly_statement": "reports",
}

def test_task_names_are_stable():
    registered = set(app.tasks.keys())
    missing = set(EXPECTED_ROUTES) - registered
    assert not missing, f"renamed or unregistered tasks: {missing}"

@pytest.mark.parametrize("task_name,queue", EXPECTED_ROUTES.items())
def test_routing(task_name, queue):
    route = app.amqp.router.route({}, task_name, args=(), kwargs={})
    assert route["queue"].name == queue

When a rename is intentional, register the old name as an alias for one release so queued messages still find a handler, then remove it. The routing mechanics are covered in Celery task routing with task_routes; this test just keeps them from drifting.

Verification

A healthy suite shows the layer split in its timings and catches the regression classes from the problem statement. Re-introduce each bug and confirm a test fails:

pytest --durations=10                      # unit + contract: all well under 100 ms
pytest -m integration --durations=10       # integration: a few seconds each

# Mutation checks (temporarily):
#  - set ignore_result=True on a chord member   -> test_statement_chord_fires_body fails
#  - pass Decimal to charge.delay                 -> serialization test fails
#  - add CrmValidationError to autoretry_for      -> test_validation_error_does_not_retry fails

Gotchas & Edge Cases

The worker thread shares the test's process. Module-level fakes and monkeypatches affect the embedded worker too — convenient, but patches applied after the worker imported something may not take effect. Patch before the fixture starts or patch at the call site.

Database transactions and the worker thread. With pytest-django, the test's transaction is invisible to the worker thread's connection. Use @pytest.mark.django_db(transaction=True) for integration tests where tasks read data the test created.

Result expiry in long suites. Session-scoped backends accumulate results; set result_expires low in test config and flush between tests.

Beat is not started by the fixtures. Test schedule definitions as data (next run times) rather than waiting for beat to fire.

FAQ

Is celery.contrib.testing.worker.start_worker different from the fixture? The fixture wraps it. Use start_worker directly as a context manager when you need a worker in a non-pytest harness or with custom pool settings.

Can I use the memory broker for integration tests? It works for simple tasks but does not support the result-backend features chords need, and its semantics differ from Redis and RabbitMQ. Use the production broker type.

How do I test a task that calls another task? In unit tests, capture the inner enqueue with the recorder from Step 3 and assert on it. In integration tests, let both run and assert on the end state.

Related