2026-07-29
Operate Redis Streams worker queues through crashes and retries
Design Redis Streams consumer groups for bounded retries, crash recovery, idempotent effects, partitioning, and production operations.
Evidence boundary
This article is source-reviewed architecture guidance. Redis commands, Java-like pseudocode, configuration choices, thresholds, and failure procedures are illustrative and unexecuted. No Redis server, Java 21 process, Spring Boot 3 application, failover, benchmark, or recovery exercise was run for this article. Verify command and client compatibility against the exact Redis, Spring Data Redis, and Lettuce versions in a canary before using any template.
Links to Redis and Spring documentation identify sourced product behavior. Queue topology, retry policy, key design, metrics, and operating thresholds are architectural recommendations. They are deliberately labeled as choices rather than Redis guarantees.
Choose the delivery contract before the data type
A production work queue needs to preserve an accepted job while a worker is unavailable, expose ownership, attempts, and recovery state for work a crashed worker did not finish, and bound repeated failure through an application retry and dead-letter policy. Redis Streams consumer groups provide the delivery-state mechanisms directly. A producer appends with XADD. A group divides new entries among named consumers through XREADGROUP. Redis records delivered but unacknowledged entries in the group's pending entries list, or PEL. A worker calls XACK only after it has established a durable outcome. Operators inspect the PEL with XPENDING, and a recovery loop transfers sufficiently idle entries with XAUTOCLAIM or targeted XCLAIM.
For retained queue data, that is an at-least-once delivery design. A worker can finish an effect and crash before XACK, so another worker can receive the same job. Conversely, Redis persistence and failover choices can lose acknowledged writes or recently accepted entries under their documented failure windows. "At least once" describes the consumer protocol, not a promise that Redis can never lose data.
Redis acknowledgement is not atomic with an arbitrary external side effect. XACK removes an entry from the PEL; it cannot atomically commit a PostgreSQL update, publish an email through a provider, finish an object-store upload, or make a payment-like irreversible call. Exactly-once processing is therefore not a property of this queue. The application must make repeat execution safe.
Use Streams when recovery ownership, delivery attempts, consumer inspection, and multiple groups justify the extra state. A Redis List reliable queue remains sound when one pool needs one FIFO backlog and application-managed recovery.
Understand the stream state machine
XADD jobs * ... appends an entry and normally lets Redis assign its stream ID. The ID orders entries within that stream. It is not the business job identifier. Put a stable application-generated job_id in the fields so duplicate publication, retry, dead-letter replay, and downstream idempotency all refer to the same intent.
Create the consumer group as deployment state, not lazily inside every worker:
XGROUP CREATE work:{ocr-03}:ready ocr-workers-v2 $ MKSTREAMMKSTREAM creates an empty stream if needed. The starting ID is a migration decision. $ means the group begins after the stream's current tail, which suits a new queue created before producers start. 0-0 makes an existing backlog eligible. Redis returns BUSYGROUP when the named group already exists, so deployment code can treat that exact condition as an idempotent success while surfacing other errors. The XGROUP CREATE reference documents these behaviors.
A consumer asks for never-delivered group entries with >:
XREADGROUP GROUP ocr-workers-v2 worker-7 COUNT 8 BLOCK 2000
STREAMS work:{ocr-03}:ready >Redis distributes entries across consumers in the group rather than broadcasting every entry. Delivery creates PEL ownership. A specific ID instead of > reads that consumer's pending history. The XREADGROUP reference documents the distinction and the need for explicit acknowledgement. Do not use NOACK for durable work because it treats a read as already acknowledged.
After processing, this removes the entry from that group's PEL:
XACK work:{ocr-03}:ready ocr-workers-v2 1848139210123-0XACK does not delete the stream entry and does not affect another consumer group. It records only that this group no longer has that ID pending, as described by the XACK reference.
The mechanism yields a useful state machine:
accepted -> available -> pending (owner A) -> durable outcome -> acknowledged
| |
+---- crash, no XACK -----+
|
idle threshold passes
v
claimable -> redelivered -> pending (new owner)There is no state named "currently executing safely." Pending means delivered and not acknowledged. It does not prove that the owner is alive, dead, stalled, or still entitled to commit an effect.
This representative command trace shows the normal path and a crash path. IDs and replies are illustrative:
> XGROUP CREATE work:{ocr-03}:ready ocr-workers-v2 $ MKSTREAM
OK
> XADD work:{ocr-03}:ready * job_id 01JABC type ocr tenant_id tenant-42 schema 3
"1848139210123-0"
> XREADGROUP GROUP ocr-workers-v2 worker-7 COUNT 1 STREAMS work:{ocr-03}:ready >
1) 1) "work:{ocr-03}:ready"
2) 1) 1) "1848139210123-0"
2) 1) "job_id" 2) "01JABC" 3) "type" 4) "ocr"
> XPENDING work:{ocr-03}:ready ocr-workers-v2
1) (integer) 1
2) "1848139210123-0"
3) "1848139210123-0"
4) 1) 1) "worker-7" 2) "1"
# worker-7 crashes before durable completion and XACK
# 120001 ms elapse, exceeding the 120000 ms minimum idle time
> XAUTOCLAIM work:{ocr-03}:ready ocr-workers-v2 recovery-2 120000 0-0 COUNT 10
1) "0-0"
2) 1) 1) "1848139210123-0"
2) 1) "job_id" 2) "01JABC" 3) "type" 4) "ocr"
3) (empty array)
> XACK work:{ocr-03}:ready ocr-workers-v2 1848139210123-0
(integer) 1The trace does not imply that a worker should wait 120 seconds in every system. That value must come from the execution timeout and failure policy. The claim transfers ownership after the illustrative first worker failed to acknowledge; the final acknowledgement is valid only after recovery-2 finds or establishes the durable OCR result.
Compare a reliable List queue honestly
The reliable List pattern atomically moves an item from a ready list to a processing list with LMOVE, or waits and moves with BLMOVE. After success, the worker removes that exact item from the processing list with LREM. A monitor requeues items that have remained in processing too long. Redis documents this pattern in the LMOVE and BLMOVE references.
That queue can also provide at-least-once delivery, but the application owns more metadata. A raw list item has no server-managed consumer owner, delivery counter, idle time, group lag, or claim operation. If processing time and attempts matter, keep a separate lease record or encode an immutable job ID and maintain metadata elsewhere. Duplicate list values make LREM especially awkward, so remove by unique job envelope rather than a repeated payload.
Fairness also differs. For one List key, Redis generally gives blocked consumers priority in blocking order when work arrives, but a client loses its position after unblocking and joins the waiting order again on its next call, as the BLPOP priority rules specify. That is not an end-to-end fairness guarantee. Streams distribute newly delivered entries across consumers in a group, but assignment is pull-driven: fast polling and processing consumers can receive more work. Neither model promises round robin.
Streams fit the OCR, notification, export, and reconciliation pools here because operators need pending ownership and crash recovery. Lists fit a compact single-pool queue without replay or multiple groups. Streams retain acknowledged entries for inspection, which also creates a retention problem.
Keep the message small and the result durable
An illustrative entry contains routing and identity, not a document or result:
job_id 01J...
type ocr
tenant_id tenant-42
subject_ref document:8f2c...
schema 3
created_at 2026-07-29T11:42:19Z
trace_id 4bf9...
attempt 0Store large OCR inputs, export files, notification templates, and reconciliation datasets in the system that owns them. The queue should carry a reference plus a schema version. Smaller payloads reduce Redis memory, replication traffic, persistence work, recovery transfer, and accidental exposure through diagnostics. Do not include credentials, full personal records, or unbounded exception text.
Recommended keys for partition 03 are:
| Purpose | Key or name |
|---|---|
| Ready stream | work:{ocr-03}:ready |
| Consumer group | ocr-workers-v2 |
| Dead-letter stream | work:{ocr-03}:dlq |
| Delayed retry schedule | work:{ocr-03}:retry |
| Optional retry payload hash | work:{ocr-03}:retry-data |
| Consumer | ocr-a17-pod-6f9d |
The braces are a Redis Cluster hash tag. They place the ready stream, dead-letter stream, and retry keys for one partition in the same hash slot when a Redis transaction or script must touch them together. The Redis Cluster specification requires multi-key operations to use one slot and documents hash tags. Do not use one global hash tag for every partition, since that would concentrate the queue on one shard.
The durable result belongs outside the stream. An OCR worker writes the extracted-text artifact, its checksum, and a completed job row before acknowledging. An export worker finalizes the object and records its immutable object key. A reconciliation worker commits the resulting database mutations and a processed-job marker. A notification worker records provider request identity and outcome. On redelivery, the worker reads this durable state and returns the already established result instead of repeating the effect.
Put transaction boundaries around intent
Producer correctness starts before XADD. If an API commits an export request to a relational database and then separately appends to Redis, a crash between those actions can leave durable business state with no job. Use a transactional outbox: in one database transaction, write the business change and an outbox row containing the stable job_id. A publisher reads committed outbox rows, calls XADD, and records dispatch. A crash can publish the same outbox row twice, so consumers still deduplicate by job_id.
On the worker side, put a database result and its idempotency record in the same database transaction whenever the effect is database-local. A unique constraint on (handler, job_id) can elect the first execution. The transaction should either contain both the business change and processed marker or neither. A marker committed before the business change can suppress needed work. A marker committed afterward in another transaction leaves a duplicate window.
External APIs need a durable state machine. Before sending a notification, persist the recipient reference, template version, provider idempotency key, and request state. Reconcile ambiguous timeouts through the provider's query interface. Without idempotent requests or lookup, the system cannot distinguish failure from success with a lost response. Redis cannot close that gap.
Redis locks alone are not sufficient fencing for irreversible side effects. A worker can pause beyond a lock TTL, another worker can acquire the lock, and the old worker can resume. The Redis distributed lock guidance explicitly recommends fencing tokens and warns against assuming a lock lives as long as its holder. Use a monotonically increasing fence that the protected resource rejects when stale, a database version predicate, or a resource-native idempotency key. If the target cannot enforce the fence, the lock only reduces overlap; it does not make an old owner harmless.
Recover crashes without manufacturing a second owner
XPENDING stream group returns group-level pending counts, ID bounds, and counts by consumer. Its extended form reports individual IDs, owners, idle time, and delivery count, with range and idle filters. This is both an operator view and an input to targeted recovery, as documented by XPENDING.
Prefer a bounded XAUTOCLAIM loop for routine crash recovery:
XAUTOCLAIM work:{ocr-03}:ready ocr-workers-v2 recovery-2
120000 0-0 COUNT 50The command scans the group PEL from a cursor, transfers entries idle for at least the threshold, and returns the next cursor. Continue until the cursor returns 0-0, then wait before another sweep. It increments attempted delivery counts when returning full messages. Since Redis 7, the response also identifies IDs whose stream entries no longer exist and clears those stale PEL references. See XAUTOCLAIM.
Use XCLAIM when an operator or repair process already has specific IDs, such as entries owned by a terminated deployment. A claim succeeds only when the entry meets the minimum idle time, and a successful claim resets its idle time. Redis serializes competing claims so only one claimant wins for that attempt. The XCLAIM reference describes this ownership transfer.
Treat idle time as a redelivery threshold, not a lease with authority. A healthy worker performing a long export can appear idle. Set the threshold above the enforced job timeout plus operating margin. Split long work into checkpointed stages or keep authoritative progress in durable storage. Heartbeats help diagnosis but cannot stop a paused former worker from resuming. Idempotency or fencing remains mandatory.
Give every process a unique consumer name tied to a deployment instance. Reusing worker-1 across overlapping processes makes ownership diagnostics ambiguous. During graceful shutdown, mark the instance unready, stop requesting new entries, let the blocking read return through a finite timeout, and finish only work that fits inside the termination budget. A completed durable result may be acknowledged. Leave incomplete jobs pending for later claim. Do not acknowledge merely to make shutdown fast.
Bound retries and isolate poison jobs
Separate failure classes in code:
- A transient dependency failure receives a delayed retry with bounded exponential backoff and jitter.
- A permanent validation failure goes directly to the dead-letter stream.
- An ambiguous external result enters reconciliation instead of blind retry.
- A process defect or payload that repeatedly crashes workers becomes a poison job.
PEL delivery count is valuable evidence, but use an explicit application attempt in the retry envelope as the policy input. Claims and history reads can affect delivery counters, and operational repair should not silently consume a business retry budget. Record a compact reason code and the prior stream ID on every rescheduled attempt.
Streams do not provide a general delayed-delivery scheduler. One illustrative design uses a sorted set scored by the next eligible epoch time. Its scheduler script or stored function must complete key-type, argument, envelope-schema, and command preflight before its first write, then process a bounded batch. For each due retry, append to the ready stream and durably write a transition marker containing the destination stream ID before deleting the sorted-set member or retry payload. Retain an authoritative retry source record until that transition is complete. On a later invocation, detect an existing marker, return its prior destination ID, and finish only the remaining cleanup. A client-side read followed by separate removal and append can lose or duplicate a retry. In Cluster, every key in that boundary must share the partition hash tag. Redis scripting isolates script execution, while its Lua API error handling documents runtime command errors; writes performed before a runtime error are not rolled back, so preflight and recoverable ordering are required. Another illustrative design leaves transient failures pending until XAUTOCLAIM reaches the idle threshold, but that couples backoff to crash recovery and can grow the PEL.
When a worker moves a transient failure from the source stream into the retry schedule, use the same recoverable ordering: complete key-type, argument, envelope-schema, and command preflight before the first write; append the retry destination and write a durable marker with its destination ID; then delete source retry state or XACK the source entry. Retain an authoritative source record until completion. A retry invocation that finds the marker returns the prior destination ID and completes remaining cleanup or XACK; it does not append a second retry entry.
After the maximum attempt, use the same recoverable ordering for the dead-letter transition. A same-slot Lua script on supported Redis versions, or a stored Redis Function on Redis 7.0 or later, must preflight every key type, argument, envelope schema, and command before its first write; process only a bounded batch; append the dead-letter record and durably write its transition marker before deleting source retry state or XACKing the source stream entry. Retain an authoritative source record until the transition completes. On retry, detect an existing marker, return the prior dead-letter ID, and finish remaining source cleanup or XACK. This is recoverable at-least-once transition logic, not rollback or exactly-once behavior: Redis scripts and Functions isolate execution, but do not undo earlier writes after a runtime error. Keep every key in that script or function in one hash slot. MULTI/EXEC makes queued writes atomic, but by itself it does not implement conditional deduplication from a read result. The record should contain job_id, tenant, handler, schema, source stream ID, attempt, reason code, and safe diagnostic references. Never acknowledge until destination append and marker durability are confirmed.
Bound crash-loop poison redeliveries separately from the application retry envelope. A process can crash after Redis delivers a message but before application code increments its retry attempt. XAUTOCLAIM increments a PEL entry's delivery count when it claims and returns that entry. Inspect those per-entry counters with the extended XPENDING form, or persist a durable crash-attempt counter, and quarantine through the same atomic dead-letter transition after an independent threshold.
A dead-letter stream is a work queue for humans and repair automation, not a graveyard. Alert on arrival, assign ownership, retain enough source data to diagnose, and define expiry. Replay creates a new stream entry that references the dead-letter ID while retaining the same business job_id unless an operator explicitly creates a new intent. Fix or quarantine poison payloads before replaying a batch.
Preserve ordering only where it is real
Stream IDs order appends within one stream. A consumer group does not guarantee completion order because workers run concurrently, jobs have different durations, failures retry later, and claims can transfer old entries. If notification B must never precede notification A for one account, do not infer that rule from stream order alone.
Partition by the ordering key and enforce the sequence at the durable target. For example, hash tenant_id + export_id to a fixed partition, then compare a per-aggregate version in the database before applying. One active worker per partition gives simpler execution order but limits throughput and still needs recovery logic. A version check handles duplicate and late completion more safely than worker count.
A single stream key is also a hot-stream boundary. All commands for that key resolve to one Redis Cluster slot. Adding consumers can increase processing capacity until Redis command traffic on that shard, network, or the downstream dependency becomes the bottleneck. It cannot distribute one stream key across masters. Create a fixed number of streams, such as 32 or 128, based on measured load and failure-domain needs. Keep partition count independent from worker count so deployments can scale without remapping every job.
Scaling workers without backpressure can turn a Redis backlog into a database or provider outage. Bound per-worker concurrency, COUNT, payload bytes, and processing time. Admit producer traffic against tenant quotas and backlog age. Pause or shed low-priority production before memory is exhausted. Scale on the oldest available job age, group lag, pending age, and service time, not XLEN alone. XLEN includes retained acknowledged history and is not queue depth.
Retain enough history to recover
Streams persist entries until explicit deletion or trimming, regardless of XACK. XTRIM MAXLEN caps approximate or exact length; MINID removes entries below an ID threshold. Approximate trimming can reduce server work. The XTRIM reference documents the strategies and warns through its reference policies that stream entries and PEL references are related but distinct.
Never choose retention solely from normal throughput. Retain entries longer than the maximum processing time, retry horizon, deployment rollback window, outage recovery objective, and investigation window. One group's oldest pending ID is not a safe trim boundary when a stream has several groups. Each group has its own pending work and first entry not yet delivered.
Redis 8.2 adds the version-specific ACKED reference policy, which lets XTRIM delete only entries acknowledged by every consumer group while still honoring the MAXLEN or MINID strategy. It is one option, not a substitute for a retention horizon: an inactive group can block reclamation indefinitely, and acknowledged history still needs an explicit time, capacity, rollback, and investigation policy. Verify server and client support before adopting it. See XTRIM and XINFO GROUPS.
For a conservative MINID procedure, coordinate group creation and XGROUP SETID with the trimmer. Enumerate every group with XINFO GROUPS, recording its last-delivered-id and lag status. For each group, get its lowest pending ID with the extended XPENDING form, and locate its first existing undelivered entry with XRANGE stream (last-delivered-id + COUNT 1. Take the earliest ID among the policy cutoff and all of those per-group boundaries, then use that ID as the exact MINID cutoff. Because MINID removes only lower IDs, every boundary remains, including each group's delivery position and the entries needed for its lag metadata. After trimming, repeat the group inventory, verify each protected pending ID still has a payload with XRANGE id id, and verify each recorded first-undelivered entry still exists. Abort the procedure if group topology changed. This is a coordinated retention operation, not a casual length cap. Default trimming can otherwise leave a PEL reference whose stream payload is gone; XAUTOCLAIM may later clear that reference, but it cannot recover the work.
Model accepted bytes, retention, stream overhead, replicas, AOF growth, and retry headroom per partition. Apply approximate MAXLEN at XADD only after reviewing the loss boundary. Use time-oriented MINID only when generated IDs and the cutoff procedure suit the policy. Monitor memory before eviction acts. A queue Redis should not use an eviction policy that can discard work.
Illustrative Java 21 and Spring Boot 3 architecture
This design is illustrative and unexecuted. Use a client API that can issue a bounded blocking XREADGROUP and manual XACK; its exact types and connection plumbing are version-specific. Spring Data documents Consumer.from, ReadOffset.lastConsumed(), manual StreamOperations.acknowledge, and the distinction between receive and receiveAutoAck in its Redis Streams reference. Use manual acknowledgement.
Separate the process into four components:
OutboxPublisherserializes a versioned minimal envelope and appends it withRedisTemplate.opsForStream().add.WorkerIntakeruns one or more bounded manual consumer-group poll loops per assigned partition, reading only new entries with>.JobHandlerestablishes the idempotent durable result and returns a classified outcome.PendingRecoveryruns a leader-elected, boundedXAUTOCLAIMscan per partition through a client API verified for the deployed Redis version.
Concise Java-like worker pseudocode:
// Illustrative and unexecuted. Types hide version-specific client plumbing.
final Semaphore admission = new Semaphore(MAX_IN_FLIGHT);
final ExecutorService offloaded = Executors.newVirtualThreadPerTaskExecutor();
final AtomicBoolean polling = new AtomicBoolean(true);
void pollLoop() {
while (polling.get()) {
try {
admission.acquire(); // Reserve capacity before Redis can deliver.
} catch (InterruptedException stop) {
Thread.currentThread().interrupt();
return;
}
boolean taskOwnsPermit = false;
try {
if (!polling.get()) {
return;
}
StreamRecord message = streams.readGroup(
GROUP, CONSUMER, readyStream, ">", count: 1, block: POLL_TIMEOUT);
if (message == null) {
continue; // The finally block returns the unused permit.
}
try {
offloaded.submit(() -> {
try {
processOffloaded(message);
} finally {
admission.release();
}
});
taskOwnsPermit = true;
} catch (RejectedExecutionException stopping) {
// This delivered entry remains in the PEL for normal recovery.
}
} finally {
if (!taskOwnsPermit) {
admission.release();
}
}
}
}
void processOffloaded(StreamRecord message) {
try {
Job job = decoder.validate(message, supportedSchemas);
try {
DurableResult result = results.find(job.jobId());
if (result == null) {
result = handler.executeIdempotently(job);
results.requireDurable(result);
}
streams.acknowledge(message.stream(), GROUP, message.id());
} catch (AmbiguousExternalResult failure) {
reconciliation.schedule(job, message.id());
// Keep pending until reconciliation establishes a durable outcome.
} catch (TransientJobException failure) {
RetryDecision decision = retryPolicy.decide(job.attempt(), failure);
if (decision.exhausted()) {
deadLetters.deadLetterAndAckAtomically(
message, job, job.attempt(), safeReason(failure));
} else {
retries.scheduleAndAckAtomically(
message, job, decision.nextAttempt(), safeReason(failure));
}
}
} catch (PermanentJobException failure) {
// This path needs only the raw record, so invalid envelopes are safe.
deadLetters.deadLetterAndAckAtomically(message, safeReason(failure));
}
}Each loop acquires exactly one permit before its COUNT 1 bounded blocking XREADGROUP, so Redis never delivers an entry without already reserved local capacity. An empty read returns that permit. A delivered entry transfers its permit only to the offloaded task; rejected submission returns the permit but deliberately leaves that entry pending in Redis for normal PEL recovery. The polling thread never runs a handler or acknowledges a record. This removes callback pause/resume coordination and its lost-wakeup race while preserving each delivery through either an offloaded task or later Redis recovery. Start multiple such bounded poll loops only when one loop's Redis round-trip latency would underuse available permits; every loop shares the same semaphore, uses a distinct consumer identity, and still uses COUNT 1.
job.attempt() is the immutable application retry attempt carried in the envelope, not a PEL delivery count. The atomic retry and dead-letter transition primitives use durable markers: on a stale-owner duplicate they find the existing marker, return the already-created destination, and finish only remaining source cleanup or acknowledgement. They therefore tolerate claims and redeliveries without creating another retry or dead-letter record.
For blocking handlers, Executors.newVirtualThreadPerTaskExecutor() gives each admitted task a virtual thread, but it does not bound concurrency. Do not size a pool of virtual threads as concurrency control. JEP 444 recommends a thread-per-task model and semaphores for limiting access to scarce services. Virtual threads do not increase database connections, provider quotas, CPU, or Redis capacity.
Use a finite poll timeout so shutdown can stop intake. First mark the instance unready, set polling false, and interrupt or otherwise wake the poll loops; wait for those loops to exit before closing the offloaded executor. Then allow only work within the termination budget to finish, leaving any incomplete entries pending for later claim. Set command, connect, and topology refresh behavior explicitly in the selected client version.
Do not hide recovery inside the ordinary listener. A separate loop needs its own rate limit, cursor, metrics, and circuit breaker so a large stale PEL cannot flood workers after an outage. Cap claimed batch size and total recovery concurrency per partition.
Deploy Redis as queue infrastructure
Persistence defines the accepted-loss window. RDB snapshots are compact and useful for backups, but writes after the last snapshot can be lost. AOF records mutations and offers always, everysec, and no-fsync policies. Redis documents that the default everysec policy may lose about one second of writes in a disaster. Review the Redis persistence guide and choose from the business recovery objective, not a generic "durable" label. Back up and restore-test the configured persistence files.
Replication and automated failover improve availability, but asynchronous replication can still lose recent accepted jobs or acknowledgements during promotion. A lost job needs an authoritative producer outbox to republish it. A lost acknowledgement causes duplicate delivery, which idempotency must absorb. Test both. If queue loss is unacceptable, Redis must not be the sole system of record for job intent.
In Redis Cluster, spread partitions across slots. Multi-key scripts, transactions, retry moves, and dead-letter moves must stay inside one hash tag. Redis Function libraries must be loaded and version-verified on every primary: Cluster does not automatically propagate function loads between primaries. See Functions in Cluster. Cluster-aware clients must handle MOVED and ASK. Reserve capacity for a promoted replica to absorb writes and recovery traffic.
Isolate tenants in the envelope, authorization path, quotas, metrics, and durable result lookup. For stronger noisy-neighbor or access boundaries, use separate partition sets, ACL users, or Redis deployments. Redis Cluster has only database 0, so logical database numbers are not tenant isolation. Redis ACLs can restrict commands and key patterns, but application authorization still must verify that a job's tenant_id matches the referenced object.
Keep Redis on a private network reachable only by trusted application and operator identities. Use TLS for client links and cluster or replication links as appropriate, unique credentials, credential rotation, and least-privilege ACLs scoped to the required stream and retry commands. Redis provides ACL and TLS guidance. Do not grant workers administrative commands, broad key access, or the ability to flush data.
Observe work, not just Redis uptime
Export these metrics per queue, partition, group, handler, and bounded tenant class:
| Metric | Meaning and action |
|---|---|
| Accepted jobs and bytes | Producer demand and payload growth |
| Group lag | XINFO GROUPS reports entries-added - entries-read, a heuristic for entries not yet delivered to the group. It can be null after arbitrary group IDs or when deleted entries invalidate metadata. |
| Oldest available age | User-visible waiting time before first delivery |
| PEL count | Delivered work without acknowledgement |
| Oldest pending idle time | Stalled work or jobs exceeding their execution budget |
| Claim count and rate | Crash recovery volume |
| Delivery attempt distribution | Retry pressure and poison jobs |
| Processing duration | Capacity, timeout, and claim-threshold input |
| Outcome count by safe reason | Success, transient, permanent, ambiguous, and dead-letter paths |
| Dead-letter arrivals and age | Unresolved operator work |
| Duplicate suppression count | Expected at-least-once behavior and abnormal redelivery spikes |
| Trimmed or missing pending payloads | Retention failure requiring immediate investigation |
| Redis memory, latency, AOF, replication, and failover state | Queue infrastructure health |
Alert on age and inability to recover, not only counts. A small queue containing one two-hour-old export can matter more than a large queue draining within its objective. Avoid tenant_id, job_id, consumer name, stream ID, or exception message as unbounded metric labels. Put those identifiers in sampled structured logs and traces. Log acceptance, delivery, durable-result decision, acknowledgement, claim, retry scheduling, and dead-letter transition with the same job_id and trace context.
XINFO GROUPS obtains lag by subtracting the group's logical entries-read counter from the stream's entries-added counter. Redis derives entries-read heuristically from last-delivered-id; it is not a scan of undelivered records. Lag is null when a group is created or set to an arbitrary ID without valid entries-read metadata, and when deleted entries create a gap between last-delivered-id and the stream's documented last-generated-id boundary. The XINFO GROUPS reference documents these cases. When lag is unavailable or suspect, use oldest available age, PEL count and age, and an application backlog derived from authoritative accepted intents minus durable outcomes.
Test the failure contract
Run focused integration and failure tests against the exact server and client versions before rollout:
- Verify group bootstrap from an empty stream,
$behavior,0-0migration, and concurrent deploys. - Kill a worker after delivery, during durable processing, after result commit, and immediately before and after
XACK. Confirm eventual completion and no duplicate durable effect. - Pause a worker beyond the claim threshold while it later resumes. Confirm idempotency or fencing rejects a stale effect.
- Exercise transient, permanent, poison, and ambiguous external failures through the retry and dead-letter limits.
- Fill each partition, enforce producer backpressure, and verify recovery does not overwhelm downstream dependencies.
- Trim near the oldest pending entry and prove the configured policy cannot remove required payloads.
- Fail over Redis with writes and acknowledgements in flight. Reconcile outbox intent, duplicate results, pending entries, and acknowledged history.
- Reshard Cluster slots and rotate credentials and certificates while workers block, reconnect, claim, and shut down.
- Restore persistence into an isolated environment and reconcile every accepted durable intent against Redis and durable results.
| Failure | Expected queue state | Required application response |
|---|---|---|
| Worker dies before result commit | Entry remains pending | Claim after threshold and execute idempotently |
Worker dies after result commit, before XACK |
Entry remains pending | Read durable result and acknowledge without repeating the effect |
XACK succeeds but response is lost |
PEL may already be clear | Treat a zero later acknowledgement as reconciliation input, not proof of lost work |
Producer crashes after XADD, before marking outbox sent |
Duplicate publication is possible | Republish and deduplicate by job_id |
Redis promotion loses a recent XADD |
Job absent from queue | Republish from the authoritative outbox |
Redis promotion loses a recent XACK |
Completed job can reappear pending | Suppress duplicate through durable result state |
| Job exceeds idle threshold while still running | Another consumer can claim it | Enforce timeout and idempotency or fencing |
| Pending payload is trimmed | PEL can reference missing data | Alert, recover from authoritative intent if possible, and fix retention |
| Poison job crashes every owner | Delivery attempts rise | Quarantine to dead-letter after a bounded policy |
| External call times out ambiguously | Side effect status is unknown | Reconcile by provider idempotency key or query before retry |
Production-readiness checklist
- The authoritative source of job intent survives Redis loss, or the accepted loss window is explicit.
- Every job has a stable
job_id, schema version, tenant identity, and minimal bounded payload. - Group creation start IDs are reviewed as part of deployment and migration.
- Workers use manual acknowledgement only after a durable result or atomic dead-letter transition.
- Database effects and idempotency records share one transaction where possible.
- External effects use resource-enforced idempotency, fencing, or a documented reconciliation path.
- Claim idle time exceeds the enforced execution budget and is not treated as proof of death.
- Retry count, backoff, maximum attempts, poison handling, and dead-letter ownership are explicit.
- Graceful shutdown stops intake, finishes bounded work, and leaves incomplete entries pending.
- Partitioning reflects ordering keys and spreads hot streams across Cluster slots.
- Concurrency and producer admission protect databases and external providers.
- Retention exceeds processing, retry, outage, rollback, and investigation windows.
- PEL age, group lag, claims, duplicates, dead letters, trim loss, persistence, and replication are alerted.
- Redis ACLs, TLS, private networking, payload controls, tenant quotas, backups, and restore tests are in place.
- Crash, pause, failover, resharding, credential rotation, and restore exercises pass on the deployed versions.
The operational takeaway is narrow: acknowledge only an outcome that can survive redelivery, and size every claim, retry, partition, and retention decision around that rule.
Sources
- Redis Streams overview
- XADD command
- XGROUP CREATE command
- XREADGROUP command
- XACK command
- XPENDING command
- XAUTOCLAIM command
- XCLAIM command
- XTRIM command
- XINFO GROUPS command
- Redis Lists overview
- LMOVE command and reliable queue pattern
- BLMOVE command
- BLPOP command
- Redis scripting
- Redis Lua API error handling
- Redis Functions
- Redis Functions in Cluster
- Redis transactions
- Redis sorted sets
- Redis distributed lock guidance
- Redis persistence guide
- Redis Cluster specification
- Redis ACL guide
- Redis TLS guide
- Spring Data Redis Streams reference
- JEP 444: Virtual Threads