Redis Streams give you an append-only log inside a single Redis key: producers append entries, readers scan ranges or block for new ones, and consumer groups split the work among workers with acknowledgements and redelivery. They are often the right tool when you already run Redis and need a durable-enough queue or event log without operating a separate broker.
Using them well depends on knowing what the structure really is. This article explains the entry-ID scheme, the memory layout that makes streams compact, the exact state a consumer group keeps, how unacknowledged entries are recovered, how trimming interacts with consumers, what asynchronous replication means for durability and how to scale past one shard. It ends with a working Python consumer, a list of failure modes and a comparison with Kafka. The basics of choosing streams over lists are covered in Redis data structures in depth; this page goes underneath them.
The model: one key, one ordered log
A stream is a value stored at a key, like a hash or a list. Each entry has an ID and a small set of field-value pairs. Entries are only appended at the end; they can be read by ID range with XRANGE and XREVRANGE, tailed with XREAD BLOCK, or distributed through a consumer group with XREADGROUP. Reading never removes an entry. Entries leave the stream only through trimming or an explicit delete.
That separates streams from the two other Redis messaging tools. A list used as a queue loses an item the moment it is popped, so a worker crash after the pop loses the work unless you build a second list for in-flight items. Pub/Sub keeps nothing: a subscriber that is disconnected misses messages forever. A stream keeps the history and tracks, per group, what was delivered and what was acknowledged.
Entry IDs: time plus a sequence
An entry ID has the form <milliseconds>-<sequence>, for example 1727740800123-0. With the automatic ID *, Redis uses the current Unix time in milliseconds and a sequence number that starts at 0 and increments for entries added in the same millisecond. IDs must strictly increase. If the server clock moves backwards, Redis keeps the last ID's millisecond part and keeps incrementing the sequence, so the ordering invariant holds even though the ID no longer matches wall-clock time exactly.
Two consequences are useful. First, IDs double as timestamps, so XRANGE orders 1727740800000 1727740860000 returns one minute of entries without any secondary index. Second, the ID gives every entry a total order within the stream, which is the only ordering guarantee you get. Explicit IDs are allowed, but each must exceed the stream's last ID or XADD fails.
Inside the key: a radix tree of listpacks
Storing every entry as a separate object would be expensive. Instead, a stream is a radix tree keyed by entry ID whose leaves are listpacks, compact contiguous byte arrays each holding many entries. The first entry of each node is a master entry; later entries store their ID as a delta from it, and when an entry has the same field names as the master entry only its values are stored. A stream of uniform events, such as orders with the same fields, therefore costs little more than its values.
Two configuration directives cap node size: stream-node-max-bytes (4096 bytes by default) and stream-node-max-entries (100 by default). Larger nodes are more compact but make deletes and range seeks inside a node slower. The layout explains several behaviours that otherwise look odd. XDEL only marks an entry deleted; the memory returns when the whole node becomes empty. Approximate trimming with ~ is cheap because it drops whole nodes rather than splitting one. And a stream with wildly varying field names compresses poorly, so keep schemas stable and put variable data in a value.
Consumer groups and the pending entries list
A consumer group is state attached to the stream: a name, a last-delivered-id, a read counter, a set of named consumers and a pending entries list, the PEL. When a consumer calls XREADGROUP GROUP billing worker-7 COUNT 100 STREAMS orders >, the special ID > asks for entries never delivered to anyone in the group. Redis hands out entries after last-delivered-id, advances it, and records each delivered ID in the PEL with its owner, its delivery time and a delivery count of one.
Each entry is delivered to one consumer in the group; different groups each see every entry, which is how one stream feeds billing and analytics independently. Calling XREADGROUP with an ID other than >, typically 0, returns the caller's own pending entries instead of new ones. A restarting worker should do exactly that first, so it finishes what it was given before it crashed.
XACK removes an ID from the PEL. Until then the entry stays pending, and the group delivers at least once: a worker that processes an entry and crashes before acknowledging will see it again, or another worker will. Handlers must therefore be idempotent, for example by recording the entry ID in the same transaction as the side effect.
Recovering stuck entries: claim, count, dead-letter
If a consumer dies, its pending entries sit in the PEL and their idle time grows. XAUTOCLAIM key group consumer min-idle-time start COUNT n (Redis 6.2 and later) scans the PEL from start and transfers entries idle longer than min-idle-time milliseconds to the calling consumer, returning them along with a cursor for the next call. The default count is 100. Claiming resets the idle time, so two workers cannot claim the same entry at once, and it increments the delivery count unless JUSTID is used. Since 7.0 the reply has a third element listing pending IDs whose entries were already trimmed or deleted; those references are removed from the PEL.
The delivery count is the poison-message detector. An entry that crashes every worker that touches it will be claimed again and again. Read counts with XPENDING in its extended form and, above a threshold, copy the entry to a dead-letter stream and acknowledge the original. Dead-letter queue architecture covers what to store with the quarantined message and how to redrive it safely.
import os, socket
import redis
r = redis.Redis(decode_responses=True)
STREAM, GROUP, DLQ = "orders", "billing", "orders:dlq"
ME = f"{socket.gethostname()}-{os.getpid()}"
MAX_DELIVERIES, IDLE_MS = 5, 60_000
def ensure_group():
try:
r.xgroup_create(STREAM, GROUP, id="0", mkstream=True)
except redis.ResponseError as e:
if "BUSYGROUP" not in str(e): # group already exists: fine
raise
def handle(entry_id, fields):
... # must be idempotent: an entry can be delivered more than once
def consume(entries):
for entry_id, fields in entries:
handle(entry_id, fields)
r.xack(STREAM, GROUP, entry_id)
def quarantine_poison():
for p in r.xpending_range(STREAM, GROUP, min="-", max="+",
count=100, idle=IDLE_MS):
if p["times_delivered"] >= MAX_DELIVERIES:
mid = p["message_id"]
rows = r.xrange(STREAM, mid, mid)
pipe = r.pipeline(transaction=True)
if rows:
pipe.xadd(DLQ, {**rows[0][1], "src_id": mid},
maxlen=100_000, approximate=True)
pipe.xack(STREAM, GROUP, mid)
pipe.execute()
def run():
ensure_group()
cursor, start = "0-0", "0" # "0": my own pending entries first
while True:
reply = r.xreadgroup(GROUP, ME, {STREAM: start}, count=100, block=5000)
entries = reply[0][1] if reply else []
if start == "0" and not entries:
start = ">" # backlog drained: switch to new entries
consume(entries)
quarantine_poison()
res = r.xautoclaim(STREAM, GROUP, ME, min_idle_time=IDLE_MS,
start_id=cursor, count=50)
cursor = res[0]
consume(res[1])
Trimming and retention
A stream grows until you trim it. XADD orders MAXLEN ~ 1000000 * ... caps the length on every append; MINID ~ <id> (6.2 and later) removes entries older than an ID, which with time-based IDs means a time-based retention window. The ~ form trims only whole nodes and may leave slightly more entries than asked; = is exact and costs more. LIMIT (6.2 and later) bounds how many entries one trim may evict, so a large backlog is removed gradually rather than in one long command.
By default trimming ignores consumer groups. An entry still pending, or never read by a slow group, is removed as readily as an acknowledged one; the PEL keeps a dangling reference, and the lag of the slow group becomes unknown. Size retention so that the slowest group's worst plausible outage fits inside it. Redis 8.2 adds options to XADD and XTRIM: KEEPREF (the default), DELREF which also removes references from every group's PEL, and ACKED which only removes entries that every group has read and acknowledged. The same release adds XACKDEL, which acknowledges and deletes in one command; with ACKED it deletes only when no group still needs the entry. Redis 8.6 adds IDMP and IDMPAUTO to XADD for idempotent producers that should not append duplicates on retry. Check your server version before relying on any of these.
Durability and replication
A stream is as durable as the Redis instance holding it. With RDB snapshots only, a crash loses everything since the last snapshot. With the append-only file and appendfsync everysec, a crash can lose roughly the last second of writes. Replication is asynchronous: the primary acknowledges XADD before replicas have it, so a failover can promote a replica that lacks the newest entries. Producers saw success and the entries are gone.
WAIT numreplicas timeout after a write blocks until that many replicas acknowledge it, which narrows the window but does not make Redis a consensus system; a failover can still pick a replica that missed the write. Consumer-group state is replicated the same way, so after a failover a group can move backwards and redeliver entries already processed. That is another reason handlers must be idempotent. If you need acknowledged writes to survive any single failure by construction, use a replicated log built for it, such as Kafka with acks=all and a minimum in-sync replica count.
Scaling: one stream is one shard
A key lives in one hash slot on one shard and every command on it runs on that shard's main thread. One stream therefore has a single-shard ceiling for both throughput and memory, no matter how large the cluster is. To scale, partition explicitly: write to orders:{0} through orders:{N-1} chosen by a hash of the order key, give each partition its own group, and assign partitions to workers. Ordering then holds per partition only, exactly as in Kafka, and changing N later reshuffles keys, so pick it with growth in mind. Rebalancing partitions across workers is your code's job; Kafka consumer group rebalancing describes the protocols a broker would otherwise run for you.
Memory is the other ceiling. Estimate bytes per entry by loading a sample and calling MEMORY USAGE on the key, then multiply by the retention length and the number of partitions on each shard, and leave room for replication buffers and fork copy-on-write during persistence.
Observing a stream
XINFO GROUPS key reports, per group, the consumer count, the PEL size, the last delivered ID, entries-read and lag, the number of entries not yet delivered to the group. Lag is computed from two counters and is reported as null after a group is created or reset at an arbitrary ID, or when entries between the group's position and the stream's end were deleted or trimmed; it recovers once the group catches up. Alert on lag growth, PEL size, oldest pending idle time and the maximum delivery count. A rising PEL with flat lag means workers receive entries and fail to acknowledge them, which is a different incident from slow workers. When producers outrun consumers for long, apply the techniques in Backpressure architecture rather than letting trimming discard unread work.
Failure modes
- Unbounded growth. No trimming, memory fills, eviction or out-of-memory errors follow. Always set MAXLEN or MINID on XADD or run XTRIM on a timer.
- Trimmed before processed. A slow group falls behind the retention window and silently loses entries. Size retention for the slowest group and alert on lag.
- Orphaned pending entries. Workers restart with new names and nobody claims the old ones. Use stable consumer names or run XAUTOCLAIM, and delete dead consumers with XGROUP DELCONSUMER after their PEL is empty.
- Poison loops. One bad entry is claimed forever. Count deliveries and dead-letter.
- Lost acknowledged writes on failover. Asynchronous replication. Make producers able to replay from their source of truth, or use a consensus-replicated log.
- Hot shard. One stream carries everything. Partition across keys.
Redis Streams or Kafka?
| Concern | Redis Streams | Kafka |
|---|---|---|
| Storage | Memory, persisted by RDB/AOF | Disk segments, optional tiered storage |
| Retention | Bounded by RAM; trim by length or ID | Days to forever |
| Replication | Asynchronous | ISR with acks=all |
| Partitioning | Manual, one key per partition | Built in, with rebalancing |
| Per-entry acknowledgement | Yes, PEL per group | No, offsets per partition |
| Operational cost | Low if Redis already runs | A cluster to operate |
Streams win for short-retention work queues, fan-out to a few groups and low-latency event passing next to an existing Redis. Kafka wins for long retention, replay of large histories, strict durability and very high throughput.
What to do next
- Run a local Redis, create a stream with XADD and inspect it with XINFO STREAM and MEMORY USAGE.
- Create a group, read with two consumers, kill one and recover its entries with XAUTOCLAIM.
- Adapt the Python consumer above, with an idempotent handler and a dead-letter stream.
- Choose a retention rule from the slowest group's worst outage and set MAXLEN or MINID on every XADD.
- Add alerts on lag, PEL size, oldest pending idle time and delivery count.
- Decide in writing what happens to acknowledged writes on failover, and test it.
- If one stream approaches a shard's limits, partition across keys before it becomes urgent.