SkillByAIOpen interactive version →

Lesson 15 / 25

Recovering Stuck Messages and Dead Letters

Claim entries from failed consumers and handle poison messages.

When a consumer never comes back

If a consumer crashes and never restarts (a pod replaced with a new name, say), its pending entries would stay stuck forever. Other consumers can take them over. XAUTOCLAIM (Redis 6.2+) scans the PEL for entries idle longer than a threshold, transfers ownership to the calling consumer, increments their delivery count and returns them for processing, using a cursor to continue across calls. XCLAIM does the same for specific IDs. Run a claim step periodically in each consumer, for example every 30 seconds with a minimum idle time comfortably above normal processing time. Use the delivery count to detect poison messages: if an entry has been delivered, say, five times and still fails, copy it to a separate dead-letter stream with error details, then XACK it in the main group so it stops blocking. Clean up consumers that no longer exist with XGROUP DELCONSUMER after their pending entries are claimed.

Claiming idle entries and dead-lettering poison messages

Entries idle for over 60 seconds are taken over; after five deliveries they are parked.

MAX_DELIVERIES = 5

def reclaim():
    cursor = "0-0"
    while True:
        # Redis 7+ returns (next_cursor, claimed_entries, deleted_ids)
        cursor, claimed, _deleted = r.xautoclaim(STREAM, GROUP, CONSUMER,
                                                min_idle_time=60_000, start_id=cursor, count=50)
        for entry_id, fields in claimed:
            info = r.xpending_range(STREAM, GROUP, min=entry_id, max=entry_id, count=1)
            deliveries = info[0]["times_delivered"] if info else 1
            if deliveries > MAX_DELIVERIES:
                r.xadd("orders:dead", {**fields, "original_id": entry_id, "reason": "max deliveries"})
                r.xack(STREAM, GROUP, entry_id)       # remove from the main PEL
                continue
            try:
                handle(entry_id, fields)
                r.xack(STREAM, GROUP, entry_id)
            except Exception:
                log.exception("retry failed %s", entry_id)
        if cursor == "0-0":
            break

Choose the idle threshold carefully

If min idle time is shorter than a slow but healthy job, two consumers will process the same entry concurrently. Set it well above your p99 processing time, and keep handlers idempotent anyway.

Quick check: What does XAUTOCLAIM do?

  • Deletes all pending entries
  • Creates a new consumer group
  • Transfers pending entries that have been idle too long to the calling consumer
  • Trims the stream
Answer

Transfers pending entries that have been idle too long to the calling consumer — XAUTOCLAIM takes ownership of idle pending entries so another consumer can process them.