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.pytestorpytest_plugins = ("celery.contrib.pytest",)inconftest.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).
freezegunortime-machinefor 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.
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.
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.
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
- Testing Background Jobs — the layered strategy this guide implements.
- Celery Chains, Groups, and Chords — the workflows the integration layer protects.
- Celery Task Retry and Error Handling — the retry settings under test.
- Chaos Testing Worker Crashes and Redelivery — the next step beyond embedded workers.