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.
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.
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.
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
- Message Size Limits & Serialization — formats and size limits.
- Compressing Job Payloads with Zstandard — the same reader-first rollout for encoding changes.
- Migrating from Redis to a Postgres Job Queue — payload compatibility during a backend move.
- Scheduled & Delayed Jobs — the long-lived payloads that set the window.