Redis Streams vs Redis Lists for Job Queues
Redis offers two data structures that can carry jobs — lists and streams — and most frameworks chose one of them long ago (Sidekiq and RQ use lists; newer systems increasingly use streams). This guide compares them directly for teams building on Redis without a framework, or choosing between frameworks, as part of Message Broker Comparison in Queue Fundamentals & Architecture. The central question is what happens to a job when a worker crashes holding it.
Problem Statement
A team runs a hand-rolled queue on Redis: producers LPUSH JSON jobs onto a list and workers BRPOP them. It is fast and simple, and it loses jobs. When a worker pod is OOM-killed or rescheduled, whatever job it had popped is gone — Redis already removed it from the list. Over a month the team estimates 0.05% of jobs vanished, including some billing reconciliations. They want at-least-once delivery with crash recovery, the ability to see which jobs are in progress and for how long, and a sense of the memory and throughput cost of the alternatives before rewriting the queue.
Prerequisites
- Redis 6.2+ (for
LMOVE/BLMOVEandXAUTOCLAIM); Redis 7 recommended. - A client library exposing list and stream commands (
redis-py,ioredis,go-redis). - Idempotent job handlers, since any at-least-once design redelivers.
- An estimate of peak backlog size, for memory planning.
Step 1 — See Why BRPOP Loses Jobs
BRPOP removes the element and returns it in one step. The job then exists only in the worker's memory; if the process dies before finishing, there is no record anywhere that it existed.
# The lossy pattern
while True:
_, raw = r.brpop("jobs", timeout=5) or (None, None)
if raw:
process(json.loads(raw)) # crash here = job gone forever
This is at-most-once delivery. It is acceptable for work that may be lost (cache warming, analytics pings) and wrong for anything else.
Step 2 — Make Lists Reliable with BLMOVE and a Reaper
The classic fix keeps lists but moves each job atomically into a per-worker processing list, removing it only after success. A reaper returns jobs from processing lists of dead workers.
WORKER = f"worker:{socket.gethostname()}:{os.getpid()}"
PROCESSING = f"processing:{WORKER}"
def work_loop():
r.set(f"heartbeat:{WORKER}", 1, ex=30)
while True:
raw = r.blmove("jobs", PROCESSING, timeout=5, src="RIGHT", dest="LEFT")
r.set(f"heartbeat:{WORKER}", 1, ex=30)
if raw is None:
continue
process(json.loads(raw))
r.lrem(PROCESSING, 1, raw) # ack: remove exactly this job
def reap():
for key in r.scan_iter("processing:*"):
worker = key.split(":", 1)[1]
if r.exists(f"heartbeat:{worker}"):
continue # worker alive
while (raw := r.rpoplpush(key, "jobs")): # return its jobs to the queue
metrics.jobs_recovered.inc()
This is the pattern behind Sidekiq Pro's reliable fetch and similar implementations. It works, but you now own heartbeats, per-worker keys, the reaper, and edge cases like duplicate identical payloads (LREM removes by value, so identical jobs need unique ids in the payload).
Step 3 — Use Streams with Consumer Groups
Streams build this bookkeeping into Redis. XADD appends an entry with a unique id; a consumer group tracks which entries have been delivered to which consumer in a pending entries list (PEL); XACK removes an entry from the PEL; XAUTOCLAIM transfers entries that have been pending too long to another consumer.
STREAM, GROUP = "jobs:stream", "workers"
r.xgroup_create(STREAM, GROUP, id="0", mkstream=True) # once; ignore BUSYGROUP error
# producer
r.xadd(STREAM, {"kind": "reconcile", "payload": json.dumps(job)}, maxlen=1_000_000, approximate=True)
# consumer
def stream_loop(consumer: str):
while True:
# 1. reclaim entries idle > 60 s from crashed consumers
_, claimed, _ = r.xautoclaim(STREAM, GROUP, consumer, min_idle_time=60_000, count=10)
# 2. read new entries
resp = r.xreadgroup(GROUP, consumer, {STREAM: ">"}, count=10, block=5000) or []
entries = claimed + [e for _, es in resp for e in es]
for entry_id, fields in entries:
process(json.loads(fields["payload"]))
r.xack(STREAM, GROUP, entry_id)
The PEL also records how many times each entry was delivered (XPENDING shows the count), which gives you a dead-letter threshold for free: an entry reclaimed more than N times can be moved to a DLQ stream instead of processed again.
Step 4 — Compare Memory and Retention
Lists delete a job when it is popped; streams keep entries until they are trimmed, even after acknowledgement. That is a feature (replay, audit) and a memory cost.
# Memory per 1M small jobs (~200-byte payloads), measured with MEMORY USAGE on a sample
redis-cli MEMORY USAGE jobs # list: ~260 MB per 1M queued items
redis-cli MEMORY USAGE jobs:stream # stream: ~230-300 MB per 1M entries (listpack-encoded)
Per entry, the two are comparable. The difference is lifetime: a list holds only the backlog, while a stream holds everything not yet trimmed. Trim with MAXLEN ~ on XADD (bounded size) or MINID (bounded age), and make sure the trim bound is well above your largest expected backlog — trimming removes entries whether or not they were acknowledged. The memory planning is covered in sizing Redis memory for queue backlogs.
Step 5 — Weigh Throughput and Operational Features
| Aspect | List + BLMOVE reaper | Stream + consumer group |
|---|---|---|
| Crash recovery | Your reaper and heartbeats | Built in: PEL + XAUTOCLAIM |
| Visibility into in-flight jobs | Scan processing:* keys | XPENDING with idle time and counts |
| Delivery count / DLQ threshold | Store in payload yourself | Built into PEL |
| Replay of past jobs | Gone after ack | Available until trimmed |
| Multiple independent consumer groups | Copy to several lists | Native: many groups per stream |
| Commands per job | ~3 (BLMOVE, LREM, heartbeats amortised) | ~3 (XADD, XREADGROUP, XACK) + periodic XAUTOCLAIM |
| Priority | Multiple lists polled in order | Multiple streams read in order |
| Framework support | Sidekiq, RQ, Resque | Newer libraries; hand-rolled |
Throughput is similar for small payloads; both are bounded by the Redis main thread, as discussed in benchmarking Redis broker throughput. The choice is about correctness tooling: streams give crash recovery and inspection without custom code.
Step 6 — Migrate the Hand-Rolled Queue to Streams
Migrate by dual-running: producers write to the stream, stream workers start, and list workers drain the remaining list backlog.
def enqueue(job: dict) -> None:
if flags.enabled("jobs_on_streams"):
r.xadd(STREAM, {"payload": json.dumps(job)}, maxlen=2_000_000, approximate=True)
else:
r.lpush("jobs", json.dumps(job))
# After the flag is fully on and LLEN jobs == 0 for a day, remove list workers.
Add a DLQ step to the stream consumer at the same time — an entry whose delivery count (from XPENDING) exceeds five is copied to jobs:dlq and acknowledged — and alert on the DLQ length.
Verification
# Kill a consumer mid-job and confirm another claims it
redis-cli XPENDING jobs:stream workers - + 10 # shows the entry, owner, idle ms
kubectl delete pod job-worker-2 --grace-period=0 --force
sleep 70; redis-cli XPENDING jobs:stream workers - + 10 # same entry, new owner, count 2
Track XLEN, PEL size (XPENDING summary count), and the oldest pending idle time in metrics. A PEL that only grows means consumers are reading but not acknowledging.
Gotchas & Edge Cases
Trimming unacknowledged entries. MAXLEN trims by position, not acknowledgement. A trim bound smaller than the backlog deletes pending jobs. Size it generously and alert when XLEN approaches it.
XREADGROUP with id "0". Reading with 0 instead of > re-reads the consumer's own pending entries — useful on startup to finish work from a previous run with the same consumer name, but easy to misuse.
Consumer name churn. Every distinct consumer name stays in the group until deleted (XGROUP DELCONSUMER). With pod names as consumer names, clean up old consumers periodically.
LREM by value. In the list pattern, two identical payloads are indistinguishable; include a unique job id in every payload.
FAQ
Which does Sidekiq use?
Lists. Open-source Sidekiq uses BRPOP semantics (with a best-effort requeue on shutdown); Sidekiq Pro's super_fetch adds reliable fetch with per-process working lists.
Are streams slower than lists?
Not meaningfully for job-queue workloads; both need a few commands per job. Streams' extra XAUTOCLAIM calls are periodic and cheap.
How do I choose the XAUTOCLAIM idle threshold?
Set it above the longest time a healthy consumer can hold an entry — p99.9 job duration plus a margin — just like a visibility timeout. Too short, and slow-but-healthy jobs are claimed by a second consumer and run twice; too long, and a crashed consumer's jobs wait that long before anyone retries them. For jobs with widely varying durations, have long-running consumers call XCLAIM on their own entries periodically (with JUSTID) to reset the idle time, which acts as a heartbeat.
Do streams survive a Redis restart?
Streams, consumer groups, and pending entries lists are persisted like any other Redis data, so they survive restarts to the extent your AOF or RDB configuration allows. With AOF everysec, up to a second of acknowledgements or new entries can be lost on a crash, which at-least-once handling absorbs.
Can I use streams for delayed jobs? No — stream entries are ordered by insertion. Keep delayed jobs in a sorted set and move them into the stream when due, as in implementing delayed jobs with Redis sorted sets.
Related
- Message Broker Comparison — Redis among the broker options.
- How to Choose Between RabbitMQ and Redis for Async Tasks — when Redis is the right broker at all.
- Redis Persistence: AOF vs RDB for Queues — durability underneath either structure.
- Deduplicating Jobs with Redis SET NX Keys — pairing at-least-once delivery with dedup.