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.

Advanced190–230 minutesAt-least-once pipeline capstoneRedis Open Source 8.10.1Free/local-firstLast reviewed: September 6, 2026

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.

01

Design the producer → Stream → consumer-group → idempotent side-effect → XACK lifecycle.

02

Demonstrate the crash-after-effect/before-ack duplicate window safely.

03

Use lag, pending count/age, consumer idle, stream length, and retry counts as backpressure/recovery signals.

04

Align retention with replay/PEL needs and distinguish producer deduplication from consumer exactly-once claims.

05

Assess persistence, replication, Cluster partitioning, security, retries, migration, rollback, and failure-injection coverage.

Exact lab baseline

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.

redis-cli · pipeline setup
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.

redis-cli · idempotent synthetic side effect
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.

redis-cli · observe the duplicate window
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.

redis-cli · claim, replay idempotently, ack
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.

redis-cli · normal successful attempt
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.

redis-cli · Stream/group operational evidence
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.

redis-cli · capstone final state
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

  1. What guarantee is the capstone designed around?
  2. Why is ACK-after-success still insufficient by itself?
  3. Which metrics distinguish not-yet-delivered work from delivered-but-unacknowledged work?
  4. Does Redis 8.6 IDMP make consumer processing exactly once?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.