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 zstandard or pyzstd, Node @mongodb-js/zstd or zstd-napi, Go github.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).

Level 3 is the sweet spot For 14 KB indexing payloads, zstd level 1 compresses 3.4 times in 45 microseconds per message, level 3 compresses 3.8 times in 70 microseconds, and level 9 compresses 4.2 times in 380 microseconds. Level 9 buys ten percent more ratio for five times the CPU. Ratio and CPU per message, 14 KB payloads level 1 3.4x, 45 µs level 3 3.8x, 70 µs: default level 9 4.2x, 380 µs Bar length is the compression ratio; the label adds compression time per message.

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).

Dictionaries capture repetition across messages A 3 KB webhook payload compressed alone shrinks to about 1.4 KB because each message has little internal repetition. With a 64 KB dictionary trained on twenty thousand real payloads, containing common field names and values, the same payload shrinks to about 600 bytes. Producers and workers must load the same dictionary version. 3 KB webhook payload raw JSON 3,000 bytes zstd alone ~1,400 bytes (2.1x) zstd + dict v1 ~600 bytes (5x) Retrain periodically as payload shapes drift; ship each version under a new id.

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.
Readers before writers Phase one deploys workers that can decode every format while producers still write plain JSON. Phase two enables compressed writes behind a per-queue flag once no old workers remain. Phase three introduces dictionary-compressed payloads with a new format tag, again only after every reader supports it. Three phases, each reversible 1. readers everywhere writers still plain JSON 2. compress per queue behind a flag 3. dictionaries new tag, same order Turning the flag off reverts writers instantly; readers stay compatible with every format.

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