skip to content

You are designing a job pipeline on Redis Streams with consumer groups. What delivery guarantee does the XREADGROUP/XACK cycle actually give, and how would you design the handlers, retry policy and failure handling around it?

level: principalimportance: should knowfreq 38%

answer

  1. at-least-once = ack after work; NOACK = at-most-once
  2. entry ID is the natural idempotency key
  3. min-idle-time ≈ retry delay, > p99 job time
  4. delivery count → threshold → your own dead-letter stream
  5. monitor pending, lag, max delivery count, stream length

basics

~20 s

At-least-once: an entry is acknowledged only after the work runs, so a crash between the two causes redelivery. Design idempotent handlers keyed on the entry ID, claim stale entries with a min-idle-time above p99 job time, cap retries via the delivery counter, and route poison entries to your own dead-letter stream.

solid answer

~60 s

The cycle is at-least-once, and only if you order it correctly: read with XREADGROUP, do the work, then XACK. A crash between the side effect and the ack means the entry stays pending and will be redelivered once a recovery sweep claims it. Using `NOACK` inverts this into at-most-once — nothing is tracked, so a failure loses the entry. That makes idempotency the handler's responsibility. The cheapest lever is the entry ID: it is unique and monotonic within the stream, so record it with the side effect (a unique key or a `SETNX processed:<id>` with a TTL longer than your retry horizon), or make the write itself naturally idempotent (upsert by business key rather than insert). Around the handler: run a periodic XAUTOCLAIM with a min-idle-time comfortably above p99 processing time so stale entries are recovered without stealing live work; treat the delivery counter as the retry count and, past a threshold, XADD the payload to a dead-letter stream and XACK the original; monitor pending count, per-consumer idle time and XINFO GROUPS lag; and trim the stream to a bound safely beyond the slowest group's pending range.

code

text · 16 lines
text
# 1. recover work this consumer already owned (after a restart)
XREADGROUP GROUP billing worker-3 COUNT 50 STREAMS orders 0

# 2. periodic sweep for entries abandoned by dead workers
XAUTOCLAIM orders billing worker-3 120000 0 COUNT 50

# 3. normal path
XREADGROUP GROUP billing worker-3 COUNT 10 BLOCK 5000 STREAMS orders >

# 4. per entry: dedupe guard, work, ack
SET processed:1712345678901-0 1 NX EX 86400   # nil => already handled, just ack
XACK orders billing 1712345678901-0

# 5. poison entry past the retry threshold
XADD deadletter:orders * src orders id 1712345678901-0 err "schema mismatch"
XACK orders billing 1712345678901-0

go deeper

for a junior

Say that the same entry can arrive twice — because the ack comes after the work — so handlers must be safe to run again.

for a middle

Explain the ordering that produces at-least-once, that NOACK gives at-most-once, and use the stream entry ID as the dedupe key.

for a senior

Give the full operational design: sweep with a justified min-idle-time, retry threshold from the delivery counter, an explicit dead-letter stream, and the pending/lag monitoring signals.

for a principal

Argue the guarantee boundary — where idempotency must live, why exactly-once is not on offer, retention bounds versus the slowest group, capacity signals, and when to shard across stream keys.

## What the mechanism actually promises The guarantee comes from where the acknowledgement sits. `XREADGROUP ... >` hands an entry to one consumer and records it in the group's pending list. The entry leaves that list only when the consumer calls XACK, or when another consumer claims it. Therefore: - If the handler completes and acks, the entry is processed once. - If the handler completes but the process dies before the ack, the entry stays pending and will be processed **again** once a claim sweep picks it up. - If the handler dies mid-work, the same thing happens — which is the point. So the end-to-end property is **at-least-once**, conditional on you actually running a recovery sweep. Nothing about the mechanism produces exactly-once: the side effect and the ack are two separate operations on two different systems and no protocol makes that pair atomic. Claiming exactly-once from Redis Streams is a red flag in an interview. The opposite choice exists explicitly: `XREADGROUP ... NOACK` skips the pending list entirely, giving at-most-once. That is a defensible choice for lossy telemetry where a lost sample is cheaper than a duplicate, but it should be argued, not defaulted into. ## Making the handler idempotent Because duplicates are structural, the design work is in the handler. **Use the entry ID as the dedupe key.** Stream IDs are unique per key and monotonically increasing, so they are a natural idempotency token. Two common shapes: - Persist the ID alongside the side effect in the same transaction of whatever store you write to — a unique constraint on it turns a duplicate into a cheap conflict. - Guard with Redis itself: `SET processed:<entry-id> 1 NX EX <ttl>` before doing the work, skipping if it already exists. The TTL must exceed the longest window in which a redelivery can occur, which is bounded by your trimming and claim policy, not by your job duration. Note this guard is itself only advisory — a crash between the guard and the side effect converts a duplicate into a loss — so prefer it where the work is expensive and the store cannot dedupe. **Or make the effect naturally idempotent.** Upsert by business key rather than insert; use conditional writes; make external calls with an idempotency token derived from the entry ID. This is strictly better than a dedupe table when the downstream system supports it. **Beware non-idempotent externals** — charging a card, sending an email. Those need the idempotency token pushed all the way into the external API, or a durable record written before the call. ## Retry policy Redis gives you two knobs and no policy engine. - **min-idle-time** on XCLAIM/XAUTOCLAIM is effectively the retry delay for a failure that killed the process. Set it above the p99 end-to-end handler time including its own internal retries and timeouts; too low and you double-process live work, too high and recovery is slow. Some designs additionally set `IDLE` on XCLAIM to push an entry's next eligible retry further out, giving crude backoff. - **the delivery counter** in the PEL is your retry count. It increments on each delivery and on each claim (unless `JUSTID` is used). For an in-process failure (the handler threw but the worker is alive) you have a choice: retry in-process with backoff while holding the entry, or leave it unacked and let the sweep bring it back later. Holding it is simpler for transient faults; releasing it distributes the retry and survives a wedge, at the cost of waiting out min-idle-time. ## Poison entries and dead-lettering Redis has no dead-letter queue. A message that crashes every handler is claimed and re-claimed forever, and in a small pipeline it can occupy enough capacity to stall progress. The pattern you must write yourself: 1. On delivery, read the entry's delivery count (from XPENDING, or track it when claiming). 2. Past a threshold — typically small, 3–5 — do not attempt the work. 3. `XADD deadletter:orders * source orders id <entry-id> payload <...> error <...>` and then `XACK orders billing <entry-id>` so it leaves the PEL. 4. Alert, because a dead-lettered entry is a defect signal, not routine. A dead-letter stream is itself a stream, so it can carry its own group for a replay tool. ## What to monitor - **Pending count** (`XPENDING` summary) — a monotonic rise means handling is slower than arrival, or something is stuck. - **Lag** (`XINFO GROUPS`, Redis 7.0+) — how many entries the group has never been delivered; the true backlog signal. - **Per-consumer idle time and pending count** (`XINFO CONSUMERS`) — distinguishes one wedged worker from a globally behind group. - **Max delivery count** across pending entries — a rising maximum is a poison entry forming. - **Stream length vs memory** — acking never trims, so retention is an independent policy. ## Retention and its interaction with recovery Trim with a capped write (`XADD ... MAXLEN ~ N *`) or `XTRIM ... MINID <id>`. The bound must exceed the worst-case pending range across all groups: trimming an entry that is still pending leaves a PEL record pointing at nothing, and the work is gone. On Redis 7.0+, XAUTOCLAIM reports and removes those dangling records, which is a cleanup, not a recovery. Choose the bound from the slowest consumer's realistic backlog, not from a memory target alone. ## Scaling shape Consumers within one group scale horizontally with no coordination, but a single stream key lives on one node (one hash slot in Cluster), so throughput ultimately bounds against that node. Partition by sharding across multiple stream keys with a deterministic assignment when you outgrow it, accepting that per-key ordering is all you keep.

  • Can you build exactly-once processing on top of this, and if so where does the guarantee actually live?
    Not from Redis alone — the side effect and the XACK are two operations that cannot be made atomic together. You get effectively-once by making the side effect idempotent: dedupe on the stream entry ID inside the destination store's own transaction, or use conditional/upsert writes. The guarantee lives in the destination, and Redis only supplies a stable unique key and redelivery.
  • How do you pick the min-idle-time for your recovery sweep?
    Anchor it above the p99 end-to-end handling time for one entry, including internal retries and external call timeouts, with headroom. Below that, sweeps steal entries from healthy but slow workers and cause genuine concurrent double-processing; far above it, work abandoned by a crashed pod waits that long before anyone resumes it. If handling times vary wildly, split the workload across streams with different sweep settings rather than picking one bad compromise.
  • The group's pending count is flat but XINFO GROUPS shows lag growing. What does that tell you?
    Entries are being written faster than consumers ask for new ones, but the entries that were delivered are being acknowledged fine — so nothing is stuck; you are simply under-provisioned or reading too slowly. Add consumers to the group, raise COUNT per read, or reduce per-entry work. A growing pending count with flat lag would be the opposite diagnosis: work is being taken and not finished.

saying these in an interview costs you the question

  • Claiming Redis Streams give exactly-once delivery
  • Acknowledging before doing the work to "avoid duplicates", which converts crashes into silent loss
  • Assuming failed entries are retried automatically without any XCLAIM/XAUTOCLAIM sweep in the design
  • Expecting a built-in dead-letter queue or a max-retries setting in Redis
  • Sizing MAXLEN purely from a memory budget, so trimming can evict entries that are still pending

context