Command Palette

Search for a command to run...

Hectal
PHASE 5Intermediate ~13 min· topic 3 of 5

Topic 5.3

Consumer Groups: Acks, Pending Entries, and Recovery

In one line

A consumer group splits a stream's entries among workers, tracks what each has received but not yet acknowledged in the Pending Entries List (PEL), and lets healthy workers claim a crashed worker's messages with XAUTOCLAIM. That gives at-least-once processing, so handlers must be idempotent.

0/5 · 0%

Think of it like this

A team sorting mail. Each letter from the pile goes to one sorter, who signs for it. When they finish, they stamp it done. If a sorter goes home sick with letters still signed out, a supervisor sees the unstamped signatures and reassigns those letters to someone else.

Key ideas

  1. 01

    Create: XGROUP CREATE orders email-svc $ MKSTREAM (start with new entries; 0 to process the whole history). Each group has its own position in the stream, so several groups (email, analytics) each get every entry, and within a group each entry goes to one consumer.

  2. 02

    Read: XREADGROUP GROUP email-svc worker-1 COUNT 10 BLOCK 5000 STREAMS orders >. The > means "entries never delivered to this group". Delivered entries enter the PEL with the consumer name, delivery time and delivery count.

  3. 03

    Acknowledge: XACK orders email-svc <id> removes the entry from the PEL after successful processing. Never acknowledged means it stays pending forever. NOACK skips the PEL (at-most-once, like Pub/Sub with history).

  4. 04

    Recover: on restart, a consumer first reads its own pending entries with ID 0 instead of >. For consumers that died for good, XAUTOCLAIM orders email-svc worker-2 60000 0-0 COUNT 100 transfers entries idle for over 60 s to worker-2 and increments their delivery count. XPENDING shows what's pending, for how long, and how many times it was delivered.

  5. 05

    Poison messages and DLQs: Redis has no built-in dead-letter queue. When the delivery count exceeds a limit (say 5), copy the entry to orders:dlq with XADD, then XACK it in the main group, so one bad message can't block retries forever.

  6. 06

    Monitoring: XINFO GROUPS orders shows pending, last-delivered-id, and (since 7.0) lag (entries not yet delivered to the group). Alert on growing lag and on old pending entries.

  7. 07

    Scaling: add consumers to the group to process in parallel; ordering across consumers is lost (each consumer sees an ordered subset). For per-key ordering, partition into several streams by key (like Kafka partitions) and give each partition one active consumer.

Code & diagrams

consumer-group.redisredis
127.0.0.1:6379> XGROUP CREATE orders email-svc $ MKSTREAM
OK
127.0.0.1:6379> XADD orders * orderId 9 status PAID
"1727520100000-0"
127.0.0.1:6379> XREADGROUP GROUP email-svc worker-1 COUNT 10 STREAMS orders >
1) 1) "orders"
   2) 1) 1) "1727520100000-0"
         2) 1) "orderId" 2) "9" 3) "status" 4) "PAID"
# worker-1 crashes before XACK ...
127.0.0.1:6379> XPENDING orders email-svc - + 10
1) 1) "1727520100000-0"
   2) "worker-1"
   3) (integer) 73012          # idle ms
   4) (integer) 1              # delivery count
127.0.0.1:6379> XAUTOCLAIM orders email-svc worker-2 60000 0-0 COUNT 100
1) "0-0"                        # next cursor (0-0 = done)
2) 1) 1) "1727520100000-0"
      2) 1) "orderId" 2) "9" 3) "status" 4) "PAID"
3) (empty array)                # IDs that no longer exist in the stream
127.0.0.1:6379> XACK orders email-svc 1727520100000-0
(integer) 1
127.0.0.1:6379> XINFO GROUPS orders
1) 1) "name"  2) "email-svc"  3) "consumers"  4) (integer) 2
   5) "pending"  6) (integer) 0  7) "last-delivered-id"  8) "1727520100000-0"
   9) "entries-read" 10) (integer) 1  11) "lag"  12) (integer) 0
worker.pypython
import redis, json

r = redis.Redis(decode_responses=True)
STREAM, GROUP, ME = "orders", "email-svc", "worker-1"
MAX_DELIVERIES = 5

def handle(fields: dict) -> None:
    # Idempotent: the provider dedupes on this key, so a redelivery is harmless.
    send_email(order_id=fields["orderId"], idempotency_key=f"email:{fields['orderId']}")

def process(entries):
    for stream, messages in entries:
        for msg_id, fields in messages:
            try:
                handle(fields)
                r.xack(STREAM, GROUP, msg_id)
            except Exception:
                pass            # stays pending; retried by XAUTOCLAIM later

def move_poison():
    for p in r.xpending_range(STREAM, GROUP, "-", "+", 100):
        if p["times_delivered"] >= MAX_DELIVERIES:
            msg = r.xrange(STREAM, p["message_id"], p["message_id"])
            if msg:
                r.xadd("orders:dlq", {**msg[0][1], "origId": p["message_id"]})
            r.xack(STREAM, GROUP, p["message_id"])

process(r.xreadgroup(GROUP, ME, {STREAM: "0"}))          # 1. my own pending first
while True:
    _, claimed, _ = r.xautoclaim(STREAM, GROUP, ME, 60_000, "0-0", count=50)
    if claimed:
        process([(STREAM, claimed)])                     # 2. rescue stuck entries
    move_poison()
    process(r.xreadgroup(GROUP, ME, {STREAM: ">"}, count=50, block=5000))  # 3. new
consumer-group.mermaiddiagram
Rendering diagram…

Interview problem

The problem

Order event processing

The order service writes events to a Redis Stream; an email consumer group and an analytics consumer group both process them. Why Streams instead of Pub/Sub? What happens if a consumer crashes? How do retries work, what happens to pending messages, how do you avoid duplicate side effects, and how do you scale consumers?

You're given

  • 3K order events/sec
  • Emails must not be lost or duplicated
  • Analytics can tolerate small delays

The interviewer follows up

01

What's the difference between XCLAIM and XAUTOCLAIM?

02

Can two consumers in a group process the same entry at the same time?

When it breaks

Consumers never call XACK (a code path returns early)

What you see

The PEL grows without limit, memory grows, and XAUTOCLAIM keeps redelivering entries that were already processed, causing duplicate work.

Fix & prevent

Acknowledge in a finally path for handled outcomes (success or moved to DLQ); alert on PEL size and oldest pending age.

Trimming with MAXLEN while a group lags far behind

What you see

Entries are trimmed before the lagging group reads them: silent data loss for that consumer. Pending entries whose data was trimmed show up as deleted IDs in XAUTOCLAIM.

Fix & prevent

Size retention to the slowest group's worst-case lag; monitor group lag; trim by time (MINID) with generous margins.

Explain it without notes

01

Explain the Pending Entries List and its role in at-least-once delivery.

02

How do you implement a dead-letter queue with Streams?

Practice

01

Create a group, read an entry with one consumer, don't ack it, then claim it with another consumer and ack it. Watch XPENDING at each step.

Trade-offs

  • ↔

    Consumer groups give load balancing and recovery, but ordering is only per consumer, and duplicates are possible.

  • ↔

    Longer claim thresholds reduce duplicate processing by slow consumers but delay recovery from crashed ones.

Done when you can

  • I can create groups, read with > and 0, ack, claim and inspect pending entries.

  • I can implement retries with delivery counts and a DLQ.

  • I design idempotent handlers for at-least-once delivery.