Compressing Job Payloads with Zstandard
Large job payloads cost broker memory, network bandwidth, and — on managed queues — money, and compression is the cheapest way to cut all three before reaching for architectural changes. This guide applies Zstandard (zstd) to job payloads as part of Message Size Limits & Serialization in Queue Fundamentals & Architecture, including trained dictionaries that make compression effective even on the small JSON messages most jobs carry.
Problem Statement
A document-indexing pipeline pushes jobs through Redis with JSON payloads averaging 14 KB (extracted text plus metadata), and a backlog after an outage reached 6 million jobs — 84 GB of Redis memory, forcing an emergency resize. A parallel SQS pipeline for webhook fan-out sends 3 KB JSON messages, and a handful exceed 256 KB, failing to enqueue. You want payload size cut substantially without changing job logic, with a format that old and new workers can both read during a rolling deploy, a CPU cost measured rather than guessed, and a clear rule for when compression is not the right fix.
Prerequisites
- Python
zstandardorpyzstd, Node@mongodb-js/zstdorzstd-napi, Gogithub.com/klauspost/compress/zstd. - A sample of real payloads (a few thousand) for measuring ratios and training dictionaries.
- Control over the serializer used by producers and workers (framework serializer hooks or a wrapper around job arguments).
- Metrics for payload size at enqueue (or the ability to add them).
Step 1 — Measure Before Compressing
Compression ratio depends entirely on the data. Measure on a real sample at a few levels, alongside the CPU time per message.
# measure.py
import json, time, zstandard as zstd, statistics
samples = [json.dumps(p).encode() for p in load_sample_payloads(5000)]
raw_sizes = [len(s) for s in samples]
print("raw mean", statistics.mean(raw_sizes))
for level in (1, 3, 9):
c = zstd.ZstdCompressor(level=level)
t0 = time.perf_counter(); out = [c.compress(s) for s in samples]; t = time.perf_counter() - t0
ratio = sum(raw_sizes) / sum(map(len, out))
print(f"level {level}: ratio {ratio:.1f}x, {t / len(samples) * 1e6:.0f} µs/msg")
# indexing jobs (14 KB text): level 1: 3.4x 45 µs | level 3: 3.8x 70 µs | level 9: 4.2x 380 µs
# webhook jobs (3 KB JSON): level 3: 2.1x 18 µs
Level 3 (the default) is usually the right trade: most of the ratio at a fraction of level 9's CPU. For large text payloads, 3–4× is typical; for small JSON, plain compression gets 1.5–2×, which dictionaries improve (Step 3).
Step 2 — Wrap Payloads in a Versioned Envelope
Never compress in place without a marker. During a rolling deploy, workers must be able to read both compressed and uncompressed payloads, and later perhaps a different dictionary. A one-byte header (or a field in the framework's headers) makes the format self-describing.
# envelope.py
import json, zstandard as zstd
FORMAT_JSON = b"\x00" # uncompressed JSON (legacy)
FORMAT_ZSTD = b"\x01" # zstd, no dictionary
FORMAT_ZSTD_DICT = b"\x02" # zstd with dictionary; next byte = dictionary id
MIN_BYTES = 1024 # below this, compression rarely pays off
_c = zstd.ZstdCompressor(level=3)
_d = zstd.ZstdDecompressor()
def encode(obj) -> bytes:
raw = json.dumps(obj, separators=(",", ":")).encode()
if len(raw) < MIN_BYTES:
return FORMAT_JSON + raw
return FORMAT_ZSTD + _c.compress(raw)
def decode(data: bytes):
tag, body = data[:1], data[1:]
if tag == FORMAT_JSON:
return json.loads(body)
if tag == FORMAT_ZSTD:
return json.loads(_d.decompress(body))
if tag == FORMAT_ZSTD_DICT:
return json.loads(decompressor_for(body[0]).decompress(body[1:]))
return json.loads(data) # no tag: pre-envelope legacy payload
The last branch keeps payloads written before the envelope existed readable, because legacy JSON starts with { or [, never with \x00–\x02. Deploy readers (workers that understand the envelope) everywhere before any producer starts writing compressed payloads.
Step 3 — Train a Dictionary for Small Messages
Small messages compress poorly because there is little repetition within one message. A dictionary trained on a sample captures the repetition across messages — field names, common values, URL prefixes — and is shared by producer and consumer.
# train.py — run offline; ship the dictionary with the code
import zstandard as zstd
samples = [json.dumps(p, separators=(",", ":")).encode() for p in load_sample_payloads(20000, kind="webhook")]
dict_data = zstd.train_dictionary(dict_size=64 * 1024, samples=samples)
open("dicts/webhook-v1.zdict", "wb").write(dict_data.as_bytes())
# runtime
DICTS = {1: zstd.ZstdCompressionDict(open("dicts/webhook-v1.zdict", "rb").read())}
_cd = zstd.ZstdCompressor(level=3, dict_data=DICTS[1])
def encode_with_dict(obj) -> bytes:
raw = json.dumps(obj, separators=(",", ":")).encode()
return FORMAT_ZSTD_DICT + bytes([1]) + _cd.compress(raw)
def decompressor_for(dict_id: int) -> zstd.ZstdDecompressor:
return zstd.ZstdDecompressor(dict_data=DICTS[dict_id])
On the 3 KB webhook payloads, the dictionary raises the ratio from 2.1× to about 5×. Dictionaries are versioned by id; never change the contents of an existing id, and keep old dictionaries loaded until no payload using them can still be in a queue (including DLQs and scheduled jobs).
Step 4 — Plug It into the Framework
Most frameworks let you register a serializer or compression codec so job code keeps passing plain objects.
# Celery: register a compression method; kombu applies it to message bodies
from kombu import compression
import zstandard as zstd
compression.register(
lambda body: zstd.ZstdCompressor(level=3).compress(body),
lambda body: zstd.ZstdDecompressor().decompress(body),
content_type="application/zstd",
aliases=["zstd"],
)
app.conf.task_compression = "zstd" # applied to task messages; recorded in headers
// BullMQ: compress data at the edges; job.data holds a base64 string
import { compress, decompress } from "@mongodb-js/zstd";
export async function addCompressed(queue: Queue, name: string, data: unknown) {
const raw = Buffer.from(JSON.stringify(data));
if (raw.length < 1024) return queue.add(name, { v: 0, d: data });
return queue.add(name, { v: 1, z: (await compress(raw, 3)).toString("base64") });
}
export async function readData<T>(job: Job): Promise<T> {
if (job.data.v === 1) return JSON.parse((await decompress(Buffer.from(job.data.z, "base64"))).toString());
return job.data.d ?? job.data; // legacy uncompressed jobs
}
Celery records the compression in message headers, so workers decode automatically — but only workers that have the codec registered. Register it in worker code first, deploy, then enable task_compression on producers. For BullMQ, base64 adds 33% on top of the compressed size; still a large net win for big payloads.
Step 5 — Account for CPU and Know the Limits
Compression moves cost from memory and network to CPU. At 70 µs per message for compression and roughly a third of that for decompression, 2,000 jobs per second costs about 0.2 CPU cores across producers and workers — negligible next to most job handlers. Measure it anyway:
# Payload bytes before and after, from metrics emitted in encode()
sum(rate(job_payload_bytes_raw_total[5m])) / sum(rate(job_payload_bytes_encoded_total[5m]))
# Redis memory per queued job after rollout
redis_memory_used_bytes / sum(queue_depth)
Compression is the wrong fix when payloads are large because they carry data that belongs elsewhere — files, full documents, large result sets. A 2 MB PDF compressed to 1.8 MB is still too big for a queue. Store it in object storage and pass a reference, as in the claim-check pattern for large payloads. Compression suits payloads that are big because of verbose structure; the claim check suits payloads that are big because of content.
Step 6 — Roll Out Safely
Order matters because producers and consumers deploy at different times:
1. Deploy readers: every worker understands the envelope (decode handles all formats).
2. Wait until no old worker versions are running (and none can be rolled back to).
3. Enable compression on producers behind a flag, one queue at a time.
4. Watch decode errors, payload size, and handler latency for a day per queue.
5. Only then introduce dictionaries (new format tag), following the same order.
A rollback of workers after step 3 would leave compressed payloads that old workers cannot read. Keep the flag so producers can switch back to plain JSON instantly, and keep readers backwards compatible permanently — the cost is a few lines of code.
Verification
def test_round_trip_all_formats():
obj = {"doc_id": "d-1", "text": "lorem " * 500}
assert decode(encode(obj)) == obj # zstd path
small = {"id": 1}
assert decode(encode(small)) == small # below threshold: plain JSON
assert decode(json.dumps(obj).encode()) == obj # legacy untagged payload
assert decode(encode_with_dict({"event": "x"})) == {"event": "x"}
In staging, enqueue the same backlog of 100,000 real payloads before and after, and compare INFO memory for Redis or ApproximateNumberOfMessages × average size for SQS. Expect the memory drop to track the measured ratio.
Gotchas & Edge Cases
Compressing already-compressed data. Images, PDFs, and gzip content do not compress further; the envelope's threshold does not catch them. Skip compression for binary content types.
Decompression bombs. A malicious or corrupt payload can expand enormously. Set a maximum output size on the decompressor (max_output_size) for payloads from untrusted producers.
Debuggability. Compressed payloads are unreadable in broker UIs. Provide a small CLI (decode-job <id>) for on-call engineers.
Cross-language producers. Every language that enqueues or consumes must implement the same envelope and dictionaries. Keep the format spec in one shared document with test vectors.
FAQ
zstd, gzip, or lz4? zstd gives better ratios than gzip at much lower CPU, and near-lz4 speed at level 1. lz4 is fine when CPU is extremely tight and ratio matters less; gzip only for compatibility.
Does SQS or RabbitMQ compress for me? No. SQS bills on the bytes you send; RabbitMQ stores what you publish. Compression happens in your producers and consumers.
What about Protobuf instead? A compact binary schema format shrinks structure overhead and pairs well with compression; see optimizing JSON vs Protobuf for job payloads.
Related
- Message Size Limits & Serialization — limits across brokers and the options for staying under them.
- Claim-Check Pattern for Large Payloads — when compression is not enough.
- Versioning Job Payload Schemas — evolving payloads safely alongside the envelope.
- Sizing Redis Memory for Queue Backlogs — the memory this saves.