Versioning Job Payload Schemas

A job's payload is a contract between code that enqueued it and code that will run it — possibly days later, possibly after several deploys. This guide shows how to change that contract without breaking jobs already in queues, as part of Message Size Limits & Serialization in Queue Fundamentals & Architecture.

Problem Statement

A subscription service renamed a job argument from plan to plan_id and added a required currency field. The deploy went out at 14:00. For the next six hours, workers failed on jobs enqueued before 14:00 (KeyError: 'plan_id'), then on retries of those jobs scheduled overnight, then — two weeks later — on a batch replayed from the dead-letter queue. Meanwhile, during the rolling deploy itself, old workers received new-format jobs and failed the other way round. You want a set of rules and code patterns under which payload changes never break a job that is queued, scheduled, retrying, or sitting in a DLQ, in either direction during a deploy.

Prerequisites

  • A single place where each job's payload is constructed (a function or class per job type) and a single place where it is parsed.
  • Knowledge of the longest time a job can sit before running: queue backlog, scheduled delay, retry horizon, and DLQ retention.
  • A serialization format with optional fields (JSON, Protobuf, Avro).
  • Code review that treats payload changes like API changes.

Step 1 — Know the Compatibility Window

Two directions of compatibility matter, and they have different lifetimes:

  • Backward compatibility (new worker, old payload): needed for as long as old payloads can exist — the maximum of backlog time, scheduled delay, retry horizon, and DLQ retention. Often days to weeks.
  • Forward compatibility (old worker, new payload): needed only during a rolling deploy and until rollback is no longer possible — typically minutes to hours.
# payload-compat-window.yaml — the numbers that decide how long to keep old code paths
max_backlog_age: 6h            # alert threshold on oldest job
max_scheduled_delay: 30d       # e.g., "renewal reminder in 30 days"
max_retry_horizon: 3d          # last retry of an exhausted backoff schedule
dlq_retention: 14d             # replays can come from here
=> backward-compat window: 30d  # support old payloads for at least this long
rollback_window: 24h           # forward-compat: new payloads must not crash old workers

The scheduled-delay line is the one teams forget: a job enqueued today to run in 30 days will be parsed by code deployed a month from now.

Old payloads outlive their deploy A payload written at deploy time can still be read hours later from a backlog, days later on its final retry, two weeks later when replayed from the dead-letter queue, and thirty days later if it was a scheduled job. The backward-compatibility window must cover the longest of these. When an old payload may still be parsed backlog 6 h retries 3 d DLQ replay 14 d scheduled 30 d: the window old-payload support must cover Forward compatibility (old workers, new payloads) only needs to last through the rollback window.

Step 2 — Prefer Additive, Optional Changes

Most changes can be made purely additive: add a field with a default, never remove or rename in the same deploy, never change a field's meaning or type.

# Before
@dataclass
class ChargeArgs:
    subscription_id: int
    plan: str

# After (additive): new optional field with a default derived from old data
@dataclass
class ChargeArgs:
    subscription_id: int
    plan: str
    currency: str | None = None          # optional: old payloads lack it

def charge(args: dict) -> None:
    a = ChargeArgs(**{k: v for k, v in args.items() if k in ChargeArgs.__dataclass_fields__})
    currency = a.currency or Subscription.get(a.subscription_id).currency   # fallback for old jobs
    ...

Two details make this robust in both directions. The parser ignores unknown keys, so an old worker receiving a new payload with currency does not crash (forward compatibility). The handler supplies a default when the field is missing, so a new worker handles old payloads (backward compatibility).

Step 3 — Add an Explicit Version and Upcast

When a change cannot be additive — restructuring, splitting a field — add a version number to the payload and convert old versions to the current shape in one function, before the handler sees them. This is "upcasting".

CURRENT = 3

def upcast(payload: dict) -> dict:
    v = payload.get("_v", 1)                      # payloads without _v are version 1
    if v == 1:
        payload = {**payload, "plan_id": payload.pop("plan"), "_v": 2}
        v = 2
    if v == 2:
        sub = Subscription.get(payload["subscription_id"])
        payload = {**payload, "currency": payload.get("currency") or sub.currency, "_v": 3}
        v = 3
    if v > CURRENT:
        raise UnknownPayloadVersion(v)            # from a newer producer: retry after deploy
    return payload

@app.task(bind=True, acks_late=True, max_retries=20)
def charge_task(self, payload: dict):
    try:
        p = upcast(payload)
    except UnknownPayloadVersion:
        raise self.retry(countdown=60)            # old worker during rollout: let a new one take it
    charge(ChargeArgsV3(**p))

Upcasters chain, so each version only knows how to reach the next. Retrying on an unknown future version is the forward-compatibility escape hatch: during a rolling deploy, an old worker hands the job back rather than failing it, and a new worker picks it up. Remove an upcaster only after the compatibility window from Step 1 has passed since the last producer wrote that version.

Upcast old payloads, handle only the current shape A version 1 payload is converted by the first upcaster, renaming plan to plan_id, into version 2. The second upcaster adds currency from the subscription record, producing version 3. A version 2 payload enters the chain at the second step. The handler only ever sees version 3. A payload with a version higher than the worker knows is retried so a newer worker can take it. v1 -> v2 -> v3 -> handler v1 payload upcast 1 to 2 plan -> plan_id upcast 2 to 3 add currency handler (v3 only) v2 payload enters here Unknown future version: retry, so a newer worker handles it during the rollout.

Step 4 — Rename Fields with Expand and Contract

A rename is the most common breaking change and the easiest to do safely across three deploys:

# Deploy 1 (expand): producers write BOTH fields; workers read new, fall back to old
enqueue("charge", {"subscription_id": 7, "plan": "pro", "plan_id": "pro"})
plan_id = payload.get("plan_id") or payload["plan"]

# Deploy 2 (switch): producers write only plan_id; workers still accept plan
enqueue("charge", {"subscription_id": 7, "plan_id": "pro"})

# Deploy 3 (contract): after the compatibility window, workers drop the fallback
plan_id = payload["plan_id"]

Each deploy is compatible with the one before and after it, so rolling deploys and rollbacks are safe at every step.

Rename in three compatible deploys Deploy one expands: producers write both plan and plan_id, and workers read plan_id with plan as a fallback. Deploy two switches producers to plan_id only while workers keep the fallback. Deploy three, after the compatibility window, removes the fallback. Any two adjacent deploys can run side by side safely. plan -> plan_id without breaking anything 1. expand write both fields read new, fall back to old 2. switch write plan_id only keep the fallback 3. contract after the 30-day window remove the fallback Put the contract date in the pull request for deploy 2, so it is not forgotten or rushed.

Deploy 3 waits for the compatibility window, not for the next sprint.

Step 5 — Validate Payloads at the Boundary

Validation at enqueue catches bad payloads before they sit in a queue for a day; validation at the worker (after upcasting) produces a clear, permanent error rather than a KeyError deep in business logic.

from pydantic import BaseModel, ValidationError

class ChargeArgsV3(BaseModel):
    model_config = {"extra": "ignore"}           # forward compatible: ignore unknown fields
    _v: int = 3
    subscription_id: int
    plan_id: str
    currency: str

def enqueue_charge(**kwargs):
    args = ChargeArgsV3(**kwargs)                # fail fast in the producer
    charge_task.delay({**args.model_dump(), "_v": CURRENT})

@app.task(bind=True, acks_late=True)
def charge_task(self, payload):
    try:
        args = ChargeArgsV3(**upcast(payload))
    except ValidationError as e:
        dead_letter(payload, reason=f"invalid payload: {e}")   # permanent: do not retry
        return
    charge(args)

A payload that fails validation after upcasting will never succeed; route it to a dead-letter queue immediately rather than burning retries, as described in Dead-Letter Queues & Poison Messages.

Step 6 — Keep Historic Payloads as Test Fixtures

Upcasters and fallbacks are code paths that production exercises rarely and at the worst times. Keep one real payload of every version as a fixture and run the current worker against all of them in CI.

# tests/test_payload_compat.py
@pytest.mark.parametrize("path", sorted(glob.glob("tests/payloads/charge/v*.json")))
def test_every_historic_version_is_accepted(path, db):
    payload = json.load(open(path))
    db.create_subscription(id=payload["subscription_id"], currency="EUR")
    assert ChargeArgsV3(**upcast(payload)).subscription_id == payload["subscription_id"]

Add a fixture whenever the version number increases, and delete it only in the same change that removes the upcaster. The broader testing approach is in Testing Background Jobs.

Verification

Before deploying a payload change, answer three questions in the pull request: can the new worker parse every payload version that could still be queued (fixtures pass)? can the old worker survive a new payload for the rollback window (unknown fields ignored, unknown versions retried)? when can the old code path be removed (date = now + compatibility window)? After deploying, watch for validation dead-letters and UnknownPayloadVersion retries; both should be zero once the rollout completes.

Gotchas & Edge Cases

Positional arguments. Frameworks that serialise positional arguments (task.delay(7, "pro")) make every change breaking. Pass a single dict or keyword arguments.

Type changes. Changing amount from integer cents to a decimal string is a new field, not a new type for the old one. Add amount_str and migrate.

Pickled payloads. Pickle ties payloads to class definitions and import paths; renaming a module breaks queued jobs. Use JSON or a schema format.

Cross-service producers. When another team's service enqueues your jobs, the payload is a public API. Publish the schema and version policy, and consider a schema registry.

FAQ

Do I need a schema registry? For jobs produced and consumed by the same codebase, versioned dataclasses and fixtures are enough. When several services produce the same payloads, a registry (Confluent Schema Registry, or a schemas repo with CI checks) enforces compatibility automatically.

How do I find out which payload versions are still in flight? Emit the payload version as a label on the enqueue and dequeue counters (jobs_dequeued_total{job="charge", payload_version="2"}). When the dequeue rate for an old version has been zero for longer than the compatibility window — and the DLQ holds none of that version — its upcaster can go. For scheduled jobs, query the scheduler's storage directly, because they will not show up in dequeue metrics until they run.

Should the version live in the payload or in headers? Either works; the payload is simplest and survives every broker and DLQ replay. Headers are cleaner when the framework supports them end to end.

What about Protobuf or Avro? Both have built-in rules for compatible evolution (new optional fields, never reuse field numbers). They help, but the queued-payload lifetime still applies — see optimizing JSON vs Protobuf for job payloads.

Related