Chapter 07 · Redis Streams and Consumer Groups
Design an At-Least-Once Stream Pipeline with Idempotent Processing, Backpressure, and Observability
Design and operate an at-least-once Redis Stream pipeline with idempotent side effects, recovery, retention, backpressure, observability, and failure-tested boundaries.
Learning outcomes
AtlasMart's fulfillment pipeline now needs a complete operating contract: append events, distribute them, make side effects idempotent, acknowledge only after success, recover stale pending work, bound memory/history, slow producers or scale consumers when backlog grows, and prove all of this with observability. The target guarantee is at-least-once processing: an event may be attempted more than once, so repeated attempts must converge on a correct application state.
Design the producer → Stream → consumer-group → idempotent side-effect → XACK lifecycle.
Demonstrate the crash-after-effect/before-ack duplicate window safely.
Use lag, pending count/age, consumer idle, stream length, and retry counts as backpressure/recovery signals.
Align retention with replay/PEL needs and distinguish producer deduplication from consumer exactly-once claims.
Assess persistence, replication, Cluster partitioning, security, retries, migration, rollback, and failure-injection coverage.
All Chapter 07 mandatory labs reuse the disposable Chapter 01
environment: Redis Open Source 8.10.1 from Docker
Official Image redis:8.10.1, container
atlasmart-redis-ch01, standalone topology, host
publication 127.0.0.1:6379, TLS disabled only
because traffic stays on loopback, default ACL user disabled,
named ACL users atlasmart-app and
academy-admin, logical database 0, AOF with
appendfsync everysec plus RDB snapshots,
persistent /data volume, and no explicit Redis
maxmemory limit or eviction policy. The primary
interface is the redis-cli shipped in the same
pinned image. Mandatory examples use only bounded synthetic
keys under atlasmart:ch07:*. No managed service,
paid broker, or external API is required.
1. End-to-end state machine
| Stage | Redis/application state | Failure question |
|---|---|---|
| Produce | XADD creates a retained Stream entry | did the producer receive the reply, and can a retry duplicate the entry? |
| Deliver | XREADGROUP assigns the entry and creates PEL ownership | did a worker receive it but crash? |
| Apply | application performs an idempotent business transition | can the same event safely run twice? |
| Acknowledge | XACK removes pending state | was ACK sent only after success? |
| Recover | XPENDING/XAUTOCLAIM transfers stale ownership | is the old worker dead or merely slow? |
| Retain | XTRIM/XADD retention bounds body history | is enough history retained for recovery/reconciliation? |
2. Create a bounded AtlasMart pipeline
Use explicit IDs so the lab can refer to events deterministically. The group starts at 0-0 to consume the seeded history.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app DEL atlasmart:ch07:pipeline atlasmart:ch07:fulfillment:1000-0 atlasmart:ch07:fulfillment:2000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XADD atlasmart:ch07:pipeline 1000-0 order 7001 action fulfilldocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XADD atlasmart:ch07:pipeline 2000-0 order 7002 action fulfilldocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XGROUP CREATE atlasmart:ch07:pipeline fulfillers 0-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XREADGROUP GROUP fulfillers worker-a COUNT 1 STREAMS atlasmart:ch07:pipeline '>'docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillers
Entry 1000-0 is now pending for worker-a.
3. Make the side effect idempotent by event identity
For the lab, fulfillment state is stored under a key derived from the Stream ID. Repeating the same HSET converges to the same state rather than incrementing a counter or creating a second shipment. Real external systems need an equivalent idempotency key or unique constraint.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HSET atlasmart:ch07:fulfillment:1000-0 order 7001 status fulfilled source-event 1000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HGETALL atlasmart:ch07:fulfillment:1000-0
This is an application modeling choice, not a Stream guarantee. If the real effect is “charge card,” the payment system must enforce the idempotency contract.
4. Failure injection: crash after effect, before XACK
Do not acknowledge 1000-0 yet. That simulates a worker that completed the side effect but lost power before XACK. The PEL correctly retains uncertainty.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillers - + 10docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HGETALL atlasmart:ch07:fulfillment:1000-0
Both facts are simultaneously true: business state says fulfilled, while Redis says the delivery attempt is still pending. This is why “exactly once” cannot be inferred from PEL/XACK mechanics.
5. Recover and re-run safely, then acknowledge
A recovery worker claims the pending event, checks/applies the same idempotent state transition, and acknowledges it only after success.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XAUTOCLAIM atlasmart:ch07:pipeline fulfillers recovery 0 0-0 COUNT 10docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HSET atlasmart:ch07:fulfillment:1000-0 order 7001 status fulfilled source-event 1000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XACK atlasmart:ch07:pipeline fulfillers 1000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillers
Repeating HSET does not create a second fulfillment record. The recovery attempt can therefore complete safely despite duplicate delivery.
6. Process the next event and observe normal flow
Now consume 2000-0, apply its idempotent effect, and acknowledge normally.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XREADGROUP GROUP fulfillers worker-b COUNT 1 STREAMS atlasmart:ch07:pipeline '>'docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HSET atlasmart:ch07:fulfillment:2000-0 order 7002 status fulfilled source-event 2000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XACK atlasmart:ch07:pipeline fulfillers 2000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillersdocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XINFO GROUPS atlasmart:ch07:pipeline
The normal path and recovery path converge on the same final state: both fulfillment records exist and pending is zero.
7. Backpressure is visible in lag, pending age, and processing latency
Backpressure means work arrives faster than the system can safely finish it. No single Redis metric captures the whole condition. Group lag measures not-yet-delivered entries when calculable; pending measures delivered-but-unacknowledged work; pending idle/delivery counts reveal stuck attempts; application processing latency and downstream saturation explain why those numbers move.
8. A practical observability panel
Collect these signals at a cadence appropriate to the workload. Do not run MONITOR or broad expensive diagnostics as your normal telemetry path.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XLEN atlasmart:ch07:pipelinedocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XINFO STREAM atlasmart:ch07:pipelinedocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XINFO GROUPS atlasmart:ch07:pipelinedocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XINFO CONSUMERS atlasmart:ch07:pipeline fulfillersdocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillersdocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app MEMORY USAGE atlasmart:ch07:pipeline
Track trends and distributions: produced entries/sec, completed entries/sec, group lag, pending count, oldest pending age, redelivery/claim rate, poison-event rate, p50/p95/p99 processing latency, Stream memory, and retention headroom.
9. Backpressure responses must address the bottleneck
If lag grows because consumers are CPU-bound, additional consumers may help. If every worker waits on the same downstream API limit, adding consumers can make the outage worse. Options include reducing producer rate, queuing at an upstream boundary, increasing safe parallelism, batching where semantics allow, partitioning work, shedding noncritical work, or fixing the downstream dependency. Scaling consumers purely from lag without saturation context is incomplete.
10. Retention must not outrun recovery
Choose MAXLEN/MINID from the maximum expected outage and
reconciliation horizon, not just average throughput. Redis 8.2
reference policies matter when trimming around PEL entries.
ACKED can protect entries until all groups
acknowledge them, but retention may then exceed a nominal MAXLEN
when outstanding references remain. Test the exact policy with
every group that matters.
11. Redis 8.6 producer idempotency solves a different duplicate source
IDMP/IDMPAUTO can deduplicate producer
retries caused by ambiguous XADD outcomes. That is useful, but
it does not remove the consumer crash window demonstrated above.
A pipeline may need producer dedupe and consumer-side
idempotency.
12. Persistence, replication, and acknowledgment are separate durability scopes
Chapter 01 uses AOF everysec plus RDB. A crash can still lose recent writes within the configured durability window; replication later in the course is asynchronous and can also have acknowledged-write uncertainty. XACK means “this group no longer needs Redis to track this delivery as pending,” not “the event and every external side effect are durably committed across all failure domains.”
13. Cluster and partitioning: one Stream key is one hot-key candidate
In Redis Cluster, one Stream key belongs to one hash slot/primary. A globally hot event Stream cannot be split across nodes without partitioning into multiple keys. Partition by tenant, region, order shard, or another stable key only after defining cross-partition ordering, consumer-group deployment, rebalancing, and observability. Multi-stream XREAD/XREADGROUP also needs current Cluster/client compatibility verification.
14. Security and tenant boundaries
ACL key patterns can constrain applications to AtlasMart prefixes, but one shared Stream is not automatically tenant isolation. Payloads can contain sensitive business data and may survive in AOF/RDB/backups beyond application retention assumptions. Use TLS on non-loopback networks, least-privilege ACLs, data minimization, backup protection, and tenant-aware application authorization.
15. Retry/timeout contract for clients
Every producer and consumer library should define connect timeout, command/block timeout, reconnect strategy, retry classification, and ambiguous-outcome handling. Retrying XADD after a lost response can duplicate production unless idempotent production is used; retrying a business effect after redelivery can duplicate side effects unless the domain operation is idempotent. Do not let generic automatic retries silently redefine delivery semantics.
16. Failure-injection matrix
| Injection | Expected Redis evidence | Application acceptance criterion |
|---|---|---|
| consumer stops after delivery | pending count/idle rises | another worker eventually claims without unsafe concurrent effects |
| crash after side effect before XACK | business state complete + entry pending | redelivery converges idempotently, then ACK |
| poison event | delivery count/retries rise | bounded retries then quarantine/reconciliation |
| producer burst | lag and/or pending trend upward | backpressure action prevents uncontrolled resource growth |
| aggressive trim | first retained ID advances | recovery policy detects whether needed history was lost |
| restart/failover later in course | persistence/replica evidence changes | documented RPO/RTO and reconciliation path hold |
17. Migration and rollback
When migrating from Lists or an external broker, dual-write/dual-read designs can create duplicates and ordering differences. Define a cutover checkpoint, event identity, dedupe horizon, reconciliation query, retention overlap, and rollback trigger. Do not delete the old queue/history until the new Stream pipeline has processed and reconciled a defined horizon under failure tests.
18. Reproducible capstone verification
Verify final Stream/group/business state and then clean up only the synthetic keys.
docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XRANGE atlasmart:ch07:pipeline - +docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XINFO GROUPS atlasmart:ch07:pipelinedocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app XPENDING atlasmart:ch07:pipeline fulfillersdocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HGETALL atlasmart:ch07:fulfillment:1000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app HGETALL atlasmart:ch07:fulfillment:2000-0docker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app MEMORY USAGE atlasmart:ch07:pipelinedocker exec -e REDISCLI_AUTH=AtlasMart-App-Lab-Only-2026 atlasmart-redis-ch01 redis-cli --user atlasmart-app DEL atlasmart:ch07:pipeline atlasmart:ch07:fulfillment:1000-0 atlasmart:ch07:fulfillment:2000-0
Acceptance criteria before cleanup: both events remain in Stream history, pending is zero, group lag is zero when calculable, and each fulfillment event has one converged state record.
19. Production judgment
A production Stream pipeline is a coordinated reliability system, not a few commands. Use stable event IDs, explicit retention, idempotent downstream effects, ACK-after-success, bounded reads, recovery with justified idle thresholds, poison-event policy, lag/pending/latency observability, and tested crash points. Define persistence and replication separately from consumer semantics; partition hot Streams deliberately; secure payloads and backups; and prove migration/rollback with reconciliation.
20. Summary and bridge to Chapter 08
You can now reason about Stream storage, replay positions, group ownership, acknowledgments, claims, retries, retention, idempotency, backpressure, and observability without claiming exactly-once processing. Chapter 08 shifts from event logs to structured document storage with Redis JSON.
Check your understanding
- What guarantee is the capstone designed around?
- Why is ACK-after-success still insufficient by itself?
- Which metrics distinguish not-yet-delivered work from delivered-but-unacknowledged work?
- Does Redis 8.6 IDMP make consumer processing exactly once?
- Why can adding consumers make backpressure worse?
Review the answers
At-least-once processing with idempotent/reconcilable side effects.
A crash after the side effect but before XACK causes redelivery, so the side effect must tolerate duplicates.
Group lag tracks not-yet-delivered work when available; pending/PEL tracks delivered-but-unacknowledged work.
No. It addresses producer retry duplicates, not consumer redelivery/side-effect semantics.
If the real bottleneck is a downstream rate limit or shared resource, more consumers increase contention rather than capacity.
Authoritative references
- Redis Streams — stream entries, consumer groups, pending entries, acknowledgments, and recovery
- XADD — entry IDs, trimming, Redis 8.2 reference policies, and Redis 8.6 idempotent production options
- XRANGE — ordered range inspection and exclusive continuation IDs
- XREAD — blocking/nonblocking reads, explicit IDs, and the special dollar ID
- XGROUP CREATE — consumer-group starting position and MKSTREAM
- XREADGROUP — group delivery, pending history, and new-message marker
- XPENDING — PEL summary/details, idle time, owner, and delivery count
- XACK — acknowledgment semantics
- XCLAIM — explicit ownership transfer and retry-count behavior
- XAUTOCLAIM — cursor-like stale-pending recovery
- XINFO GROUPS — pending count, last-delivered ID, entries-read, and lag
- XTRIM — exact/approximate retention and PEL reference policies
- Redis 8.10 release notes — pinned server release family
- Idempotent message processing — Redis 8.6 producer-side deduplication and its scope
- MEMORY USAGE — per-key memory evidence
- Redis persistence — AOF/RDB durability scope
- Redis Cluster specification — slot locality and partitioning context