Architecture

Redis as an Event Broker: Pub/Sub and Streams Patterns for Microservices Without Kafka

Ruslan Ismailov Published 12 min read
R

Introduction: When Kafka Is Overkill and Why Redis Has Become a Popular Choice in 2026

Kafka is a powerful tool, but deploying and operating it demands significant resources: a separate ZooKeeper or KRaft cluster, replication configuration, broker monitoring, and team training. For startups, small teams, or microservice systems with moderate load, these costs are often hard to justify.

In 2026, Redis is increasingly used as an event broker in event-driven architectures. There are several reasons: Redis is already present in most production stacks as a cache, it's well understood by developers, easy to operate, and provides two mature messaging mechanisms — Pub/Sub and Streams. If your team lacks dedicated infrastructure expertise and your message volumes are in the thousands rather than millions per second, Redis as a Kafka alternative deserves serious consideration.

Redis Pub/Sub: How It Works and Delivery Guarantees

How Pub/Sub Works

Redis Pub/Sub is a classic in-memory publish-subscribe model. A publisher sends a message to a channel using the PUBLISH command, and all active subscribers on that channel instantly receive it through a persistent connection. There is no buffering: if a subscriber is not connected at the time of publishing, the message is lost permanently.

Redis supports both exact channel subscriptions and pattern-based subscriptions via PSUBSCRIBE, which allows subscribing to groups of channels using a wildcard, such as orders.*.

Delivery Guarantees — Or the Lack Thereof

Pub/Sub in Redis operates on a fire-and-forget basis: no persistence, no acknowledgments, no redelivery. This is not a bug — it's a design feature. If a consumer crashes, restarts, or simply falls behind the producer, messages are lost. This semantics is called "at most once."

When Pub/Sub Is Appropriate

  • Cache invalidation: notifications about data changes where losing a single event is not critical
  • Real-time UI notifications: user online status, chat messages
  • Broadcast events without guarantees: system metrics, heartbeat messages
  • Prototyping event-driven flows before introducing a more complex solution

Redis Streams: Architecture, Consumer Groups, and the ACK Mechanism

What Are Redis Streams

Redis Streams, introduced in Redis 5.0, are a persistent, ordered data structure resembling an append-only log. Each entry has a unique ID in the format 1704067200000-0 (milliseconds + sequence number). Streams are stored in memory and can be replicated and persisted using standard RDB/AOF mechanisms.

Consumer Groups

The key difference from Pub/Sub is support for consumer groups. A consumer group processes a stream collaboratively: each message is delivered to exactly one member of the group. This enables horizontal scaling of processing — multiple workers read from the same stream without duplicating work. Different groups receive all messages independently, similar to Kafka topics with multiple consumer groups.

ACK Mechanism and Pending Entries

After receiving a message, the consumer must explicitly acknowledge processing using the XACK command. Until acknowledged, the message remains in the Pending Entries List (PEL). If a consumer crashes, unprocessed messages can be reassigned to another worker using XCLAIM. This provides "at least once" semantics.

Persistence and Retention

Streams support length limits via MAXLEN — you can retain the last N messages or trim by time. When using AOF or RDB, data survives a Redis restart. However, it's important to understand: Redis is an in-memory database, and in the event of total data loss (without replication and persistence), the stream will be gone.

Comparison: Redis Pub/Sub vs Redis Streams vs Kafka

Below are the key parameters of the three solutions that matter most when choosing an event broker for a microservice architecture.

  • Persistence: Pub/Sub — no; Streams — yes (RDB/AOF); Kafka — yes (disk log)
  • Delivery guarantees: Pub/Sub — at most once; Streams — at least once; Kafka — at least once / exactly once
  • Consumer groups: Pub/Sub — no; Streams — yes; Kafka — yes
  • Message reprocessing: Pub/Sub — not possible; Streams — yes, via XCLAIM; Kafka — yes, via offset reset
  • Message ordering: Pub/Sub — not guaranteed; Streams — guaranteed within a stream; Kafka — guaranteed within a partition
  • Operational complexity: Pub/Sub — minimal; Streams — low; Kafka — high
  • Throughput: Pub/Sub — high; Streams — high (hundreds of thousands/sec); Kafka — very high (millions/sec)
  • Long-term storage: Pub/Sub — none; Streams — limited (RAM); Kafka — virtually unlimited
  • Infrastructure dependencies: Pub/Sub — Redis only; Streams — Redis only; Kafka — Kafka + ZooKeeper/KRaft

Event-Driven Pattern with Redis Streams: Implementation in Go

Producer

An example Go producer using the go-redis/v9 library. The producer publishes an order placement event to the orders stream.

package main

import (
    "context"
    "fmt"
    "log"
    "time"

    "github.com/redis/go-redis/v9"
)

func main() {
    ctx := context.Background()

    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    defer rdb.Close()

    // Publish event to stream
    id, err := rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "orders",
        MaxLen: 10000,    // stream length limit
        Approx: true,     // approximate trimming for performance
        Values: map[string]interface{}{
            "order_id":   "ORD-12345",
            "user_id":    "USR-67890",
            "amount":     "1500.00",
            "currency":   "USD",
            "created_at": time.Now().Unix(),
        },
    }).Result()

    if err != nil {
        log.Fatalf("XAdd error: %v", err)
    }

    fmt.Printf("Event published, ID: %s\n", id)
}

Consumer with Consumer Group

The consumer creates a group (if it doesn't exist), reads messages, and acknowledges processing via XACK.

package main

import (
    "context"
    "fmt"
    "log"
    "os"
    "os/signal"
    "syscall"
    "time"

    "github.com/redis/go-redis/v9"
)

const (
    streamName = "orders"
    groupName  = "order-processor"
)

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    defer rdb.Close()

    consumerName := fmt.Sprintf("worker-%s", os.Getenv("POD_NAME"))

    // Create group if it doesn't exist
    err := rdb.XGroupCreateMkStream(ctx, streamName, groupName, "$").Err()
    if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
        log.Fatalf("XGroupCreate error: %v", err)
    }

    // Handle signals for graceful shutdown
    sigCh := make(chan os.Signal, 1)
    signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)

    go func() {
        <-sigCh
        cancel()
    }()

    log.Printf("Consumer %s started, group: %s", consumerName, groupName)

    for {
        select {
        case <-ctx.Done():
            log.Println("Shutting down gracefully")
            return
        default:
        }

        // Read up to 10 messages with a 2-second timeout
        streams, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
            Group:    groupName,
            Consumer: consumerName,
            Streams:  []string{streamName, ">"},
            Count:    10,
            Block:    2 * time.Second,
        }).Result()

        if err != nil {
            if err == redis.Nil || err.Error() == "redis: nil" {
                continue // timeout, no new messages
            }
            if ctx.Err() != nil {
                return
            }
            log.Printf("XReadGroup error: %v", err)
            time.Sleep(time.Second)
            continue
        }

        for _, stream := range streams {
            for _, msg := range stream.Messages {
                if err := processOrder(msg); err != nil {
                    log.Printf("Processing error for %s: %v", msg.ID, err)
                    // Message remains in PEL for reprocessing
                    continue
                }

                // Acknowledge successful processing
                if err := rdb.XAck(ctx, streamName, groupName, msg.ID).Err(); err != nil {
                    log.Printf("XAck error for %s: %v", msg.ID, err)
                }
            }
        }
    }
}

func processOrder(msg redis.XMessage) error {
    fmt.Printf("Processing order: %s, amount: %s\n",
        msg.Values["order_id"],
        msg.Values["amount"],
    )
    // Business logic for order processing
    return nil
}

Error Handling and Dead Letter: What to Do with Unprocessed Messages

One of the common challenges with Redis Streams is messages "stuck" in the PEL. A message can end up there if the consumer crashed before calling XACK, or if processing consistently fails.

PEL Handling Strategy

The recommended pattern is to periodically run a task that uses XPENDING to check for messages stuck longer than a threshold (e.g., 5 minutes) and moves them to another consumer via XCLAIM. After N failed attempts, the message is transferred to a separate "dead letter" stream.

func reclaimStalePendingMessages(ctx context.Context, rdb *redis.Client) {
    minIdleTime := 5 * time.Minute

    pending, err := rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream: streamName,
        Group:  groupName,
        Start:  "-",
        End:    "+",
        Count:  100,
    }).Result()

    if err != nil {
        log.Printf("XPendingExt error: %v", err)
        return
    }

    for _, p := range pending {
        if p.Idle < minIdleTime {
            continue
        }

        if p.RetryCount >= 5 {
            // Move to dead letter stream
            moveToDeadLetter(ctx, rdb, p.ID)
            rdb.XAck(ctx, streamName, groupName, p.ID)
            continue
        }

        // Reassign to current worker
        rdb.XClaim(ctx, &redis.XClaimArgs{
            Stream:   streamName,
            Group:    groupName,
            Consumer: "reclaim-worker",
            MinIdle:  minIdleTime,
            Messages: []string{p.ID},
        })
    }
}

func moveToDeadLetter(ctx context.Context, rdb *redis.Client, msgID string) {
    rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "orders:dead-letter",
        Values: map[string]interface{}{
            "original_id": msgID,
            "moved_at":    time.Now().Unix(),
        },
    })
    log.Printf("Message %s moved to dead letter stream", msgID)
}

Scaling Consumers in Kubernetes

Deploying Workers

Each worker Pod in Kubernetes uses the POD_NAME environment variable as a unique consumer identifier within the group. Redis balances the load among group members automatically — no additional coordination required.

apiVersion: apps/v1
kind: Deployment
metadata:
  name: order-processor
spec:
  replicas: 3
  selector:
    matchLabels:
      app: order-processor
  template:
    metadata:
      labels:
        app: order-processor
    spec:
      containers:
        - name: worker
          image: myapp/order-processor:latest
          env:
            - name: POD_NAME
              valueFrom:
                fieldRef:
                  fieldPath: metadata.name
            - name: REDIS_ADDR
              value: redis-service:6379
          resources:
            requests:
              cpu: 100m
              memory: 64Mi
            limits:
              cpu: 500m
              memory: 256Mi

Autoscaling Based on Stream Length

KEDA (Kubernetes Event-Driven Autoscaling) supports Redis Streams out of the box. The scaler reads the PEL length or the number of unprocessed messages and scales the Deployment automatically.

apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: order-processor-scaler
spec:
  scaleTargetRef:
    name: order-processor
  minReplicaCount: 1
  maxReplicaCount: 20
  triggers:
    - type: redis-streams
      metadata:
        address: redis-service:6379
        stream: orders
        consumerGroup: order-processor
        pendingEntriesCount: "50"  # scale if PEL > 50

Monitoring Redis as an Event Broker: Key Metrics and Tools

Key Metrics

  • Stream length (XLEN orders) — a growing length signals that consumers are falling behind
  • PEL size (XPENDING orders order-processor - + 1) — the number of unacknowledged messages
  • Group lag — the difference between the last ID in the stream and the last acknowledged ID by the group
  • used_memory — Redis stores everything in RAM; memory exhaustion is critical
  • connected_clients — the number of active consumer connections
  • evicted_keys — if Redis evicts keys due to memory pressure, data may be lost

Tools

For monitoring Redis as an event broker, use the Redis Exporter + Prometheus + Grafana stack. Redis Exporter (oliver006/redis_exporter) exports metrics such as redis_stream_length, redis_stream_groups, and redis_stream_pending_entries_count. Set up alerts for: stream length > 10,000; PEL > 500; group lag > 60 seconds; used_memory > 80% of maxmemory.

Redis Broker Limitations: When You Actually Need Kafka

Honesty matters: Redis as an event broker has fundamental limitations you should know about upfront.

  • Storage is limited by RAM. Kafka stores data on disk and can retain messages for weeks. Redis, when memory runs low, starts evicting data — and messages are lost if MAXLEN is not configured and memory is not monitored.
  • No native partitioning. One stream is one sequence. For scaling, you need to manually create multiple streams and distribute load among them.
  • Replication is asynchronous. In Redis Sentinel and Cluster, replication is asynchronous. During a failover, a few of the most recent messages may be lost.
  • No exactly-once semantics. Redis Streams provide at least once. For exactly once, idempotency must be implemented at the application level.
  • No built-in Schema Registry. Message format compatibility is the developer's responsibility.

Choose Kafka if: you need to store events for months, your volume exceeds several million messages per second, strict exactly-once semantics are required, or you have a dedicated infrastructure team to manage it.

Consider RabbitMQ if you need richer message routing (exchanges, routing keys, dead letter queues out of the box) at moderate volumes.

Real-World Case: Replacing Redis Pub/Sub with Redis Streams in Production

Context

An e-commerce platform team was using Redis Pub/Sub to send order status change notifications. A notification microservice subscribed to the order.status.changed channel and sent push notifications and emails. The system worked, but incidents occurred periodically: when a consumer Pod restarted in Kubernetes, some notifications were lost — users did not receive their delivery emails.

The Problem

Pub/Sub does not buffer messages. While a Pod was restarting (30–60 seconds), all events were lost. A rolling update with zero downtime didn't help: during the Pod switchover, there was always a brief subscription gap.

The Solution

The team migrated to Redis Streams with a consumer group called notification-service. The producer (order management service) started writing to the order-events stream instead of publishing to a channel. The consumer reads via XREADGROUP and acknowledges with XACK only after a notification is successfully sent. A PEL worker checks for stuck messages every minute and retries processing.

The Result

In the three months following the migration, the team recorded zero incidents involving lost notifications. The migration took two days, including testing. No additional infrastructure dependencies were introduced — Redis was already in the stack. Operational maintenance costs remained unchanged.

Summary and Recommendations for Choosing the Right Tool

Redis as an event broker is a pragmatic choice for teams who want event-driven architecture without the operational complexity of Kafka.

  • Use Redis Pub/Sub for broadcast notifications, cache invalidation, and real-time events where losing a message or two is not critical.
  • Use Redis Streams for reliable asynchronous task processing, inter-service event exchange, and any scenario where at-least-once delivery guarantees matter.
  • Add KEDA for automatic consumer scaling in Kubernetes based on stream length or PEL size.
  • Always configure MAXLEN on streams and monitor used_memory — this protects against data loss when RAM runs low.
  • Implement a dead letter stream and XCLAIM logic for correct failure handling.
  • If your load grows to millions of events per second, long-term storage is required, or exactly-once semantics are needed — migrate to Kafka. Redis itself will make it clear when that time comes.

Redis Streams is not a "poor man's Kafka." It's a different tool with different trade-offs — one that is a perfect fit for the majority of microservice systems running at moderate load.

Technologies

Tags

Ruslan Ismailov

Senior Web / Backend Developer. Senior web/backend developer with 9 years of experience. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservices, CI/CD. More about me →