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

Project Goals

  • Stream orders:events with 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:

bash
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 ID milliseconds-sequence
  • Fields are strings; JSON payload in one field is common in apps

Step 2: Create Consumer Group

bash
XGROUP CREATE orders:events fulfillment-workers $ MKSTREAM

$ — group reads only new messages after creation. Use 0 to consume from beginning.

Read group info:

bash
XINFO GROUPS orders:events

Step 3: Consumer Loop (redis-cli Concept)

Terminal A (consumer worker-1):

bash
XREADGROUP GROUP fulfillment-workers worker-1 COUNT 10 BLOCK 5000 STREAMS orders:events >

Code explanation:

  • > means undelivered messages for this consumer
  • BLOCK 5000 waits up to 5 seconds for new data

When messages return, process and acknowledge:

bash
XACK orders:events fulfillment-workers 1716890400000-0

Repeat XREADGROUP in a loop in application code.

Step 4: Idempotent Processing

Workers may see the same logical event twice after XCLAIM. Store processed IDs:

bash
SET processed:1716890400000-0 1 NX EX 86400

If SET NX fails, skip—already handled.

Application pattern:

text
Read message → check idempotency key → update DB → XACK

Align order state with MySQL transactions in real apps.

Step 5: Pending and Crash Recovery

List stuck messages:

bash
XPENDING orders:events fulfillment-workers

Claim idle messages (worker died):

bash
XCLAIM orders:events fulfillment-workers worker-2 60000 1716890400000-0

Code explanation:

  • 60000 — claim entries idle longer than 60 seconds
  • New worker worker-2 takes 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:

java
// Concept: listener on "orders:events", group "fulfillment-workers"
// Deserialize fields, call orderService.updateStatus(), ack message id

Details in Spring Boot Redis integration and Redis Spring chapter.

Verification Checklist

  • XADD appends events
  • Group consumer receives with XREADGROUP
  • XACK removes from pending for that consumer
  • Stopped consumer leaves XPENDING; XCLAIM recovers
  • 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.