Project: Stream Message Queue
Introduction
Redis Streams with consumer groups provide durable, acknowledged messaging—closer to a real queue than Pub/Sub or List BRPOP. This project models order status updates: producers append events, workers in a group consume with XREADGROUP, and XACK confirms processing. You learn idempotency and pending-entry recovery patterns used in production microservices.
Prerequisites
- Streams
- Pub/Sub and Lightweight Queues—when not to use Streams
- Spring Boot Integration (optional Java consumer)
Project Goals
- Stream
orders:eventswith order id, status, timestamp - Consumer group
fulfillment-workers - Producer:
XADD - Worker:
XREADGROUP, process,XACK - Handle pending messages after worker crash with
XPENDING/XCLAIM
Step 1: Define Event Shape
Each stream entry is field-value pairs:
XADD orders:events * order_id 1001 status "PAID" ts "2026-05-28T10:00:00Z"
XADD orders:events * order_id 1001 status "SHIPPED" ts "2026-05-28T14:00:00Z"Code explanation:
*auto-generates entry IDmilliseconds-sequence- Fields are strings; JSON payload in one field is common in apps
Step 2: Create Consumer Group
XGROUP CREATE orders:events fulfillment-workers $ MKSTREAM$ — group reads only new messages after creation. Use 0 to consume from beginning.
Read group info:
XINFO GROUPS orders:eventsStep 3: Consumer Loop (redis-cli Concept)
Terminal A (consumer worker-1):
XREADGROUP GROUP fulfillment-workers worker-1 COUNT 10 BLOCK 5000 STREAMS orders:events >Code explanation:
>means undelivered messages for this consumerBLOCK 5000waits up to 5 seconds for new data
When messages return, process and acknowledge:
XACK orders:events fulfillment-workers 1716890400000-0Repeat XREADGROUP in a loop in application code.
Step 4: Idempotent Processing
Workers may see the same logical event twice after XCLAIM. Store processed IDs:
SET processed:1716890400000-0 1 NX EX 86400If SET NX fails, skip—already handled.
Application pattern:
Read message → check idempotency key → update DB → XACKAlign order state with MySQL transactions in real apps.
Step 5: Pending and Crash Recovery
List stuck messages:
XPENDING orders:events fulfillment-workersClaim idle messages (worker died):
XCLAIM orders:events fulfillment-workers worker-2 60000 1716890400000-0Code explanation:
60000— claim entries idle longer than 60 seconds- New worker
worker-2takes ownership and retries
Step 6: Dead Letter Pattern (Concept)
After N failed attempts, XADD to orders:events:dlq and XACK original to stop redelivery—or trim with documented ops policy.
Step 7: Spring Boot Consumer (Sketch)
Use StreamMessageListenerContainer or Lettuce xreadgroup in @Scheduled / dedicated thread:
// Concept: listener on "orders:events", group "fulfillment-workers"
// Deserialize fields, call orderService.updateStatus(), ack message idDetails in Spring Boot Redis integration and Redis Spring chapter.
Verification Checklist
-
XADDappends events - Group consumer receives with
XREADGROUP -
XACKremoves from pending for that consumer - Stopped consumer leaves
XPENDING;XCLAIMrecovers - Idempotency prevents double ship on retry
FAQ
Stream vs Pub/Sub?
Pub/Sub drops offline subscribers—use Streams when you need persistence and acks.
Stream vs List BRPOP?
Lists are simpler FIFO without consumer groups or pending tracking.
Trim stream size?
XTRIM orders:events MAXLEN ~ 10000 caps memory.
Multiple consumers same group?
Redis load-balances messages across consumers in one group.
Exactly-once delivery?
Impossible end-to-end—aim for at-least-once + idempotent handlers.
Order total ordering?
Single stream preserves order; shard by order_id hash for scale.