Testing Background Jobs

Background jobs fail in ways request handlers never do — they are retried, redelivered, run concurrently, and interrupted mid-way — and this guide lays out a testing strategy that catches those failures before production, as part of Backend Frameworks & Worker Scaling. Most teams test the logic inside a job and stop there. The bugs that page people at night live elsewhere: a job that is not safe to run twice, an enqueue that happens outside the transaction, a retry policy that retries the unretryable, a chord that never fires because a test ran it synchronously.

The goal is not maximal coverage. It is a small set of tests at each layer that mirror the failure modes a queue introduces: at-least-once delivery, retries with backoff, concurrent workers, and crashes at the worst possible moment.

The Scenario: Green Tests, Red Production

A team's Celery suite has 400 passing tests. Every task has unit tests; CI runs them in about two minutes using task_always_eager=True. In one quarter, production suffered three job incidents: a chord whose body never ran because a member task had ignore_result=True (eager mode does not use the result backend, so the test passed); a report task that sent emails twice after worker restarts (no test ever ran a task twice); and a payment task that retried a card-declined error eight times (the retry configuration was never exercised in tests at all).

None of those bugs was in the task's logic. All three were in the interaction between the task and the queue — the layer the test suite had carefully mocked away.

Layers of job tests The base layer is many fast unit tests of handler logic with fakes. The middle layer tests contracts: that jobs are enqueued inside transactions, with the right arguments, and that retry settings classify errors correctly. The top layer is a small number of integration tests with a real broker and worker, including duplicate delivery, crash, and redelivery scenarios. Where each class of bug is caught real broker + worker chords, crashes, redelivery contract tests enqueue in txn, args, retry classification, run twice handler unit tests business logic with fakes: fast, many All three incidents in the scenario lived in the top two layers, which the suite skipped.

Architectural Overview: What a Job Test Must Control

A background job has more inputs than its arguments. To test it meaningfully, a test needs control over five things, and the layers above differ in how many they exercise:

  1. Arguments and state — the payload and the database or external state the handler reads. Every layer controls these.
  2. Delivery count — whether this is the first delivery or the third. Needed to test idempotency and retry-dependent behaviour.
  3. Time — backoff delays, scheduled jobs, TTLs, and deadlines. Tests must fake the clock rather than sleep.
  4. Failures of dependencies — timeouts, 5xx, rate limits, connection drops — injected deterministically.
  5. Process lifecycle — the worker being killed after the side effect but before the acknowledgement.
# conftest.py — the controls, gathered as fixtures
import pytest
from freezegun import freeze_time

@pytest.fixture
def clock():
    with freeze_time("2026-09-18 12:00:00") as frozen:
        yield frozen                      # clock.tick(timedelta(seconds=30))

@pytest.fixture
def flaky(monkeypatch):
    """Make a dependency fail N times, then succeed."""
    def install(target: str, failures: int, exc: Exception):
        calls = {"n": 0}
        real = _resolve(target)
        def wrapper(*a, **kw):
            calls["n"] += 1
            if calls["n"] <= failures:
                raise exc
            return real(*a, **kw)
        monkeypatch.setattr(target, wrapper)
        return calls
    return install

@pytest.fixture
def run_twice():
    """Invoke a handler twice with identical input to simulate redelivery."""
    def go(fn, *args, **kwargs):
        fn(*args, **kwargs)
        fn(*args, **kwargs)
    return go

Faking the clock matters more for jobs than for most code. A retry test that sleeps for real backoff delays either takes minutes or uses shortened delays that do not match production — and then tests a configuration nobody runs.

Implementation 1: Unit and Contract Tests Without a Broker

The cheapest valuable tests call the handler function directly and assert on side effects, then check the contract between the job and the queue: that it is enqueued at the right moment, with the right arguments, and that its retry configuration classifies errors correctly.

# test_send_invoice.py — Celery task tested as a plain function plus its contracts
from tasks import send_invoice
from celery.exceptions import Retry

def test_sends_invoice_once_when_run_twice(db, mailer, run_twice):
    inv = db.create_invoice(status="open")
    run_twice(send_invoice.run, inv.id)             # .run bypasses the broker entirely
    assert mailer.sent_count(inv.id) == 1           # idempotency, not just correctness

def test_transient_error_retries(db, flaky, monkeypatch):
    inv = db.create_invoice(status="open")
    flaky("tasks.mailer.send", failures=1, exc=ConnectionError())
    monkeypatch.setattr(send_invoice, "retry", lambda **kw: (_ for _ in ()).throw(Retry()))
    with pytest.raises(Retry):
        send_invoice.run(inv.id)

def test_permanent_error_does_not_retry(db, flaky):
    inv = db.create_invoice(status="open")
    flaky("tasks.mailer.send", failures=99, exc=InvalidRecipient("bounced"))
    send_invoice.run(inv.id)                        # handled: marked failed, no Retry raised
    assert db.invoice(inv.id).email_status == "invalid_recipient"

def test_enqueued_only_after_commit(db, captured_tasks):
    with pytest.raises(RuntimeError):
        with db.transaction():
            create_invoice_and_enqueue(db, amount=100)
            raise RuntimeError("rollback")
    assert captured_tasks("tasks.send_invoice") == []   # nothing escaped the rolled-back txn

The last test is the one most suites lack. It needs a way to capture enqueues without a broker — Celery's task_always_eager is the wrong tool (it runs the task), so patch apply_async or use your framework's test mode. Sidekiq has Sidekiq::Testing.fake!, BullMQ can be pointed at an in-memory mock, and River ships rivertest helpers. Framework-specific setups are in testing Celery tasks with pytest and testing Sidekiq jobs with RSpec.

What eager mode skips In production, a task call is serialized, sent through the broker, picked up by a worker, and its result stored in the result backend; retries go back through the broker. In eager mode the call runs inline in the test process, so serialization errors, broker routing, result-backend-dependent features like chords, and real retry scheduling are never exercised. Production path vs eager shortcut .delay() serialize broker worker result backend .delay() eager: runs inline, every box above skipped Serialization errors, routing, chords, and real retries only show up on the top path.

Implementation 2: Integration Tests Against a Real Broker

The top layer runs real workers against a real broker in a container, started by the test suite. It is slower — seconds per test rather than milliseconds — so keep it to the behaviours that genuinely need it: workflows (chains, chords, flows, batches), serialization of real payloads, routing to queues, and crash/redelivery.

# test_integration.py — Celery with a real Redis via testcontainers
import pytest
from testcontainers.redis import RedisContainer

@pytest.fixture(scope="session")
def redis_url():
    with RedisContainer("redis:7.2") as r:
        yield f"redis://{r.get_container_host_ip()}:{r.get_exposed_port(6379)}/0"

@pytest.fixture(scope="session")
def celery_config(redis_url):
    return {"broker_url": redis_url, "result_backend": redis_url,
            "task_acks_late": True, "task_serializer": "json"}

def test_statement_chord_runs_body(celery_app, celery_worker, fake_regions):
    result = statement_workflow("c-17", "2026-08").apply_async()
    assert result.get(timeout=30) is None
    assert storage.exists("pdf/c-17.pdf")        # body ran: chord counter worked

The same approach works in Node with BullMQ — see integration testing BullMQ workers with Testcontainers — and in Go with River against a Postgres container. Start one container per test session, not per test, and isolate tests with separate queue names or a flush between tests.

Trade-off Analysis: Which Layer Tests What

Behaviour Handler unit test Contract test Real broker test
Business logic Best — Redundant
Idempotency (run twice) Good Good Realistic but slow
Enqueue inside the transaction — Best Good
Retry classification Possible with mocks Best Realistic
Backoff timing With fake clock With fake clock Slow unless configurable
Serialization of real payloads — Partial Best
Chords / flows / batches — — Only option
Crash after side effect, before ack — — Only option

A reasonable budget for a service with 30 job types: every job has unit tests and a run-twice test; every enqueue site has a transaction contract test; every retry configuration has a classification test; and there are perhaps 10–20 integration tests covering each workflow and one crash scenario per critical job.

Failure Modes the Tests Must Reproduce

Duplicate delivery. Every job will eventually run twice. The run-twice test is the minimum; for jobs with external side effects, assert on the external effect (emails sent, charges made), not on return values.

Crash between side effect and acknowledgement. The worker sends the email and is killed before acking; the broker redelivers. Only a real-broker test with a killable worker reproduces this — covered in chaos testing worker crashes and redelivery.

Retries that should not happen. A permanent error retried with backoff wastes capacity and may repeat partial side effects. Test the classification table explicitly: for each exception type, assert retry or no retry.

Backoff that is wrong. A misconfigured backoff (seconds instead of milliseconds, no cap, no jitter) only shows under load. Test the delay function directly with a fixed random seed, as in testing retry and backoff logic deterministically.

Concurrent execution of the same job. Two workers pick up duplicates simultaneously. A test that runs the handler in two threads with a barrier before the side effect exposes missing locks or unique constraints.

The crash only a real-broker test can reproduce Worker A receives a job, sends the email, and is killed before acknowledging. After the visibility timeout the broker redelivers the job to worker B, which sends the email again unless the handler checks an idempotency record. The test kills worker A at exactly that point and asserts that only one email was sent. Kill after the side effect, before the ack worker A receive send email SIGKILL worker B redelivered email again? The assertion: exactly one email, whatever point the kill lands on.

Testing Concurrency: Two Workers, One Job

Duplicate deliveries do not always arrive one after the other. With at-least-once brokers, a visibility-timeout expiry or a redelivery during a rebalance can hand the same job to two workers at the same time. A handler that is idempotent when run sequentially — "check whether the email was sent, then send it" — can still send twice when both copies pass the check before either records the send.

The test for this runs the handler in two threads and forces both to reach the critical point before either proceeds. A barrier placed inside a patched dependency does it deterministically, without relying on timing luck:

import threading

def test_concurrent_duplicates_send_once(db, mailer, monkeypatch):
    inv = db.create_invoice(status="open")
    barrier = threading.Barrier(2)
    real_lookup = tasks.already_sent

    def synchronised_lookup(invoice_id):
        result = real_lookup(invoice_id)   # both threads read "not sent"...
        barrier.wait(timeout=5)            # ...then both proceed together
        return result

    monkeypatch.setattr(tasks, "already_sent", synchronised_lookup)
    threads = [threading.Thread(target=send_invoice.run, args=(inv.id,)) for _ in range(2)]
    [t.start() for t in threads]; [t.join() for t in threads]
    assert mailer.sent_count(inv.id) == 1

A check-then-act handler fails this test; the fix is to make the claim atomic — an INSERT ... ON CONFLICT DO NOTHING into a sends table before sending, a row lock, or a unique constraint that makes the second writer fail. The test then passes because exactly one thread wins the insert. Idempotent consumers with Postgres unique constraints covers the fixes.

Check-then-act fails under concurrency With a check-then-act guard, both threads read that the email has not been sent, both pass the barrier, and both send. With an atomic insert into a sends table, both threads attempt the insert, one succeeds and sends, and the other sees a conflict and returns without sending. Two copies of one job, released together check, then send T1: not sent T2: not sent both send: 2 emails atomic claim T1: insert wins T2: conflict one email The barrier makes the race deterministic, so the test fails every time the guard is not atomic.

Asserting on Metrics and Logs

The operational signals a job emits are part of its behaviour. An on-call engineer relies on job_failures_total rising when a job fails, on a structured log line carrying the job id and attempt number, and on a DLQ metric when a job is abandoned. If a refactor silently drops those, the next incident is harder to debug — and no functional test notices.

Cheap assertions close the gap. Most metrics libraries expose their registry in-process, and logging frameworks can be captured:

from prometheus_client import REGISTRY

def metric(name, **labels):
    return REGISTRY.get_sample_value(name, labels) or 0.0

def test_permanent_failure_is_counted_and_logged(db, flaky, caplog):
    inv = db.create_invoice(status="open")
    flaky("tasks.mailer.send", failures=99, exc=InvalidRecipient("bounced"))
    before = metric("job_failures_total", task="send_invoice", kind="permanent")
    send_invoice.run(inv.id)
    assert metric("job_failures_total", task="send_invoice", kind="permanent") == before + 1
    record = next(r for r in caplog.records if r.getMessage() == "job failed")
    assert record.job_id and record.attempt == 1 and record.error_kind == "permanent"

Keep these to the signals that alerts and runbooks depend on. The logging fields worth asserting on are the ones described in Structured Logging for Workers, and the metrics are those behind the alerts in SLOs & Alerting for Job Queues.

Performance Tuning the Suite

Job test suites get slow when every test starts a worker. Keep them fast:

  • Run handlers directly in unit tests (task.run(...), job.perform(...), worker.Work(ctx, job)). No broker, no serialization.
  • Share containers per session, flushing Redis or truncating queue tables between tests instead of restarting containers.
  • Configure retry delays through settings, so integration tests can run the real retry machinery with delays scaled to milliseconds while unit tests verify the production values.
  • Run the real-broker layer in parallel CI jobs, separate from unit tests, so a two-minute unit suite stays two minutes.
  • Fail fast on timeouts: integration tests should wait for conditions with a bounded poll (wait_until(predicate, timeout=10)) rather than fixed sleeps.
# .github/workflows/test.yml — split layers so feedback stays fast
jobs:
  unit:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - run: pytest -m "not integration" -n auto        # seconds
  integration:
    runs-on: ubuntu-latest
    services:
      redis: { image: "redis:7.2", ports: ["6379:6379"] }
      postgres: { image: "postgres:16", env: { POSTGRES_PASSWORD: test }, ports: ["5432:5432"] }
    steps:
      - uses: actions/checkout@v4
      - run: pytest -m integration --timeout=120        # a few minutes, runs in parallel with unit

Testing Workflows and Scheduled Jobs

Two kinds of job logic resist the approaches above and deserve a note of their own.

Multi-step workflows — chains, chords, BullMQ flows, Sidekiq batches, sagas — have behaviour that exists only between jobs: the fan-in fires once, a failure in step three triggers compensation of steps one and two, the workflow record reaches a terminal state. Test them at the integration layer with injected failures at each step boundary, and assert on the final state of the workflow and the side effects left behind. The patterns in implementing sagas with compensating jobs include a parametrised test that fails each step in turn; that shape generalises to any workflow.

Scheduled and periodic jobs add time as an input. Test the schedule definition separately from the job: a unit test that the cron expression or schedule object produces the expected next run times across a daylight-saving boundary, and an ordinary job test for what the job does when it runs. Never test a periodic job by waiting for it to fire.

# Schedule definitions are data: test them like data
from celery.schedules import crontab
from datetime import datetime
from zoneinfo import ZoneInfo

def test_nightly_reconcile_runs_once_per_day_across_dst():
    sched = crontab(hour=2, minute=30)
    tz = ZoneInfo("Europe/Berlin")
    # 2026-03-29 is the spring-forward date in Europe: 02:30 does not exist locally
    after = datetime(2026, 3, 28, 3, 0, tzinfo=tz)
    runs = next_runs(sched, after, count=3, tz=tz)
    assert [r.date().isoformat() for r in runs] == ["2026-03-29", "2026-03-30", "2026-03-31"]

The DST case in that test is exactly the kind of edge covered in timezone-safe job scheduling across DST; encoding it in a test keeps it from regressing when someone "simplifies" the scheduler configuration.

Testing Payload Compatibility Across Deploys

A rolling deploy runs old and new code side by side, so for a few minutes new producers enqueue jobs that old workers consume, and old producers enqueue jobs that new workers consume. A job whose arguments changed shape — a renamed field, a new required parameter — fails in exactly that window, and then again for any jobs that were scheduled or retrying across the deploy. Unit tests never see it because both sides of the test use the same version of the code.

A compatibility test keeps a small corpus of serialized payloads from previous releases and runs the current handler against each. When a payload change is intentional, the test forces the author to handle the old shape explicitly, usually by giving the new field a default:

# tests/fixtures/payloads/send_invoice/v1.json, v2.json ... committed as they ship
import glob, json

@pytest.mark.parametrize("path", sorted(glob.glob("tests/fixtures/payloads/send_invoice/*.json")))
def test_current_handler_accepts_historic_payloads(path, db, mailer):
    payload = json.load(open(path))
    db.create_invoice(id=payload["invoice_id"], status="open")
    send_invoice.run(**payload)                 # must not raise TypeError / KeyError
    assert mailer.sent_count(payload["invoice_id"]) == 1

Add a new fixture file whenever a job's arguments change, and delete fixtures only after the longest possible retry or schedule horizon has passed. The versioning rules that make this workable are in versioning job payload schemas.

FAQ

Is task_always_eager ever acceptable? For testing that calling code triggers a job with the right arguments and that the job's logic works in a simple flow, it is convenient. Never rely on it for anything involving chords, retries, serialization, routing, or acknowledgement — it bypasses all of them.

Should integration tests use the production broker type? Yes. Behaviour differs between Redis, RabbitMQ, SQS, and Postgres queues in exactly the places that matter — redelivery timing, ordering, visibility. LocalStack or ElasticMQ are good stand-ins for SQS in CI.

How do I test jobs that call external APIs? Record-and-replay (VCR-style cassettes) or a fake server for unit and contract tests, with injected failures for the retry cases. Keep one smoke test against a sandbox of the real API, run on a schedule rather than on every commit.

Where do flaky job tests usually come from? Almost always from real time: fixed sleep calls waiting for a worker, backoff delays measured on the wall clock, or scheduled jobs that fire "soon". Replace sleeps with bounded polling on a condition, fake the clock in unit tests, and make retry delays configurable so integration tests can shrink them without changing the code under test.

How many crash tests are enough? One per job type whose side effect is expensive or irreversible — payments, emails, external writes. Cheap, naturally idempotent jobs rarely justify them.

Related