Redis: More Than a Cache, Building Event-Driven Systems with Redis Streams

Redis: More Than a Cache, Building Event-Driven Systems with Redis Streams

# architecture# backend# database# systemdesign
Redis: More Than a Cache, Building Event-Driven Systems with Redis StreamsRahman Nugar

Yes you heard it right. Redis is more than a cache, if you've probably heard of Redis before, you may...

Yes you heard it right. Redis is more than a cache, if you've probably heard of Redis before, you may have implemented it as a distributed cache, built a distributed rate limiter on it. However, as systems grow, services need a reliable way to communicate without being tightly coupled and this is where event-driven communication patterns become useful. Redis Streams provides a medium for this event-driven communication.

Redis streams according to Redis Labs is a data type that provides a super fast in-memory abstraction of an append only log. You can think of it as a log where you can store messages/events. But doesn't Redis use memory, is that not ephemeral? why would you want to store messages in Redis? I was asked this question on X(Twitter) more recently, (I apologize in advanced for the overly expressive language).

Yes, by default, Redis operates primarily in-memory, meaning data could be lost during a crash if persistence isn't enabled. However Redis has several data structures(when implemented) that shifts away from this, this means you can store messages durably in Redis by either appending on every write to log(AOF) or taking timely snapshots(RDB snapshots), there's no perfect storage implementation so you simply use which aligns with your needs.

This article walks through the process of building event-driven systems and using Redis streams as a as a lightweight message broker, for everything Redis, you can check out their official blog

Table of Contents

  1. What are Distributed Systems?
  2. Event-Driven Communication Patterns
  3. Pub/Sub Model
  4. Why a Traditional Queue Isn't Enough
  5. Redis Streams Introduction
  6. XREADGROUP: Coordinating Consumer Groups
  7. Pending Entries List: Handling Failed Consumers
  8. Building a Redis Stream Worker in Go
  9. Redis Streams vs Kafka and RabbitMQ

1. What are Distributed Systems?

The word distributed is self explanatory however we often conflate distributed systems as microservices alone(welp). Simply put, a distributed system is a collection of independent processes or computers that communicate with each other over a network to achieve a common goal. Now, you'd ask doesn't every computer communicate over networks already? Yes, however it is not just to communicate over a network but what the communicate does(achieve some common goal, behave as a single logical system).

A monolith underneath can have multiple internal components that communicates over a network, i.e, your application might be a single running deploy, but it communicates with a separate database, cache, message queue, or external services all over a network. It also encompasses having multiple instances of your application, your database having replicas or being split into shards however it is easier to grasp the concept using microservices where each service is a an independently deployable application.

Distributed systems

2. Event-Driven Communication Patterns

We've introduced distributed systems and what they are essentially but we have to talk about how they communicate over networks. There are two major means of communication we have to talk about here; synchronous and asynchronous communication. Synchronous communication is commonly done through request/response mechanisms such as HTTP or gRPC, where the caller waits for a response. This article focuses on asynchronous, event-driven communication.

An event-driven architecture is a software design approach where systems communicate by producing and reacting to events asynchronously rather than making direct synchronous calls. If you notice the picture in section 1, the services are connected via blue lines and red lines. We can say the blue lines are event driven while the red lines are synchronous, simplistic but it shows how services communicate using either pattern in a distributed system.

We ideally use event-driven communication when we want services to be loosely coupled. Services can communicate asynchronously without directly depending on each other, meaning the producer doesn't have to wait for the consumer to finish before continuing.

3. Pubsub Model

Pubsub means publisher-subscriber, it is a model where one device, node, application publishes a message and other applications subscribed to the publisher receive this message.
Pubsub Model

An implementation of Pub/Sub using Redis is usually ephemeral; the message is lost once published, meaning subscribers that are unavailable at that moment will miss the message. Other implementations, such as Google Pub/Sub, provide durable message retention, allowing subscribers to consume messages after they become available again.

Pub/Sub is mostly useful when the publisher doesn't need to know or wait for its subscribers, but what if we need an event/job or particular piece of work to actually be picked up and processed?

4. Why a Traditional queue isn't enough

A traditional message queue works differently to a regular pubsub model. Instead of publishing a message to multiple subscribers, messages are published to a queue and a consumer picks up the message to process it (yes, a queue, your regular FIFO data structure).

Queue

A queue is effective for background jobs where a particular job needs to be processed by one consumer. Now, job queue libraries like BullMQ (Node.js), Asynq (Go), and Celery (Python) provide a more robust implementation of this pattern, handling things like retries, acknowledgements, scheduling, and worker management, so you may ask why not just use them and why do we need Redis streams?

5. Redis streams Introduction

Like I initially explained, Redis Streams is a Redis data type that provides an append-only log where you can store messages/events. Unlike regular Redis Pub/Sub, messages written to a stream are persisted according to your Redis persistence configuration and remain available to be consumed later.

In this snippet, we can see a set of sequential numbers, they are the stream IDs which Redis generates for each message/event entry. The ID consists of a timestamp and a sequence number. Redis atomically generates these IDs, so even if multiple events are added within the same millisecond, they can still have unique, ordered IDs.

The order_id and event hare are the actual data stored in the stream entry.

orders
---------------------------------------------------------
1706000000000-0   order_id=123   event=order.created
1706000001000-0   order_id=124   event=order.paid
1706000001000-1   order_id=125   event=order.updated
---------------------------------------------------------
Enter fullscreen mode Exit fullscreen mode

You are probably wondering, Okay? Why Redis streams if job queues can store events too. The key difference is that Redis Streams gives us a persistent event log that multiple independent consumer groups can consume events independently.". A job queue is mainly concerned with getting a job to a worker to be processed. With Streams, the same event can be consumed by different services independently, and each consumer group keeps track of its own progress(I'll explain what a consumer group is).

                order.created
                     │
               Redis Stream
              /      |       \
             ↓       ↓        ↓
        Inventory  Analytics  Notifications
Enter fullscreen mode Exit fullscreen mode

In this snippet, we have an event (order being created). If we were using a queue, consumers would have to compete for messages in the same queue, so essentially we'd have one consumer process the event while the others miss it. If we have multiple services or consumers that need the same event for downstream work, we would typically need separate queues for each service/consumer.

However Redis Streams allows multiple consumer groups to consume the same stream independently. Now that we understand the difference, let's talk about the fundamentals of Redis Streams.

Adding events with XADD

The first command we need to know is XADD. As the name suggests, XADD is used to add/append an entry to a stream.

XADD orders * order_id 123 event created
Enter fullscreen mode Exit fullscreen mode

Here, orders is the name of the stream, * tells Redis to generate the stream ID(i.e 1706000000000-0) for us, and order_id and event are the fields and values stored in the entry.

Adding an event seems quite simple, but operationally, we need to think about what happens when our application updates its database successfully but fails to publish the event to Redis.

If we write changes to our database and call Redis, there's a chance that the message may be lost if Redis is unavailable, to prevent against this we use an outbox pattern.

Outbox Pattern

This isn't primary to Redis streams but I just needed to quickly explain this so we understand how to add to any message broker and ensure the message isn't lost.

Instead of writing changes to the database and then calling Redis separately, we can instead have a table(usually called an outbox table) where we write to in the same transaction as our business operation. We store the business change to its own table and the event/mesage to the outbox table, notice we didn't call Redis or sent anything to Redis yet, we instead adopt a different approach to doing that.

We can have a worker poll for changes to a database table continuously, we can also use something like Listen/Notify where the database wakes the worker(isolated to Postgres) as our primary mechanism and use the worker poll as a fallback as the Listen/Notify mechanism is ephemeral and message can be lost if worker is unavailable, we can also use CDC(change data capture) tools like Debezium, where the tool reads from the database WAL(write ahead log) directly and stream it to its intended channel. Now that we have more reliable transport, we can then transport the events into Redis streams reliably.

Reading from a stream with XREAD

Now that we've added an entry to the stream, we need a way to read it. XREAD is used to read entries from a stream.

XREAD STREAMS orders 0
Enter fullscreen mode Exit fullscreen mode

Here, orders is the stream we want to read from and 0 tells Redis to start reading from the beginning of the stream.

orders
---------------------------------------------------------
1706000000000-0   order_id=123   event=order.created
1706000001000-0   order_id=124   event=order.paid
1706000001000-1   order_id=125   event=order.updated
---------------------------------------------------------
Enter fullscreen mode Exit fullscreen mode

This is essential useful when we want to read the entire stream, but for a consumer that is only interested in new events, starting from 0 would mean reading all the existing entries first.

Redis Streams also allows us to provide a specific stream ID to tell Redis where to start reading from.

XREAD STREAMS orders 1706000001000-0
Enter fullscreen mode Exit fullscreen mode

Here, 1706000001000-0 is the ID we want to start reading from. Redis will return the entries from that point onward.

1706000001000-0   order_id=124   event=order.paid
1706000001000-1   order_id=125   event=order.updated
Enter fullscreen mode Exit fullscreen mode

We can also decide to read only new entries by using $, which tells Redis to start reading from the end of the stream.

XREAD STREAMS orders $
Enter fullscreen mode Exit fullscreen mode

However, if there are no new entries when XREAD is executed, Redis immediately returns nothing and our consumer is left polling for new events, we can avoid this uneccessary polling by using blocking reads.

XREAD BLOCK 5000 STREAMS orders $
Enter fullscreen mode Exit fullscreen mode

BLOCK 5000 tells Redis to wait for up to 5 seconds for a new entry. The connection stays open while Redis waits, and Redis returns the new entry as soon as one arrives.

Now the stream doesn't have infinite storage and as the stream grows, we may no longer need some entries, so Redis allows us to delete a specific entry using XDEL.

XDEL orders 1706000001000-0
Enter fullscreen mode Exit fullscreen mode

Here, Redis deletes the entry with the specified stream ID. However, instead of deleting entries one by one, we can also limit how large the stream is with MAXLEN.

MAXLEN can either be approximate or exact; exact means Redis will keep the stream at or below the specified number of entries, while approximate means Redis may temporarily keep slightly more entries to reduce the work Redis needs to do when trimming the stream during XADD.

XADD orders MAXLEN = 10000 * order_id 126 event order.created
Enter fullscreen mode Exit fullscreen mode
XADD orders MAXLEN ~ 10000 * order_id 126 event order.created
Enter fullscreen mode Exit fullscreen mode

The = symbol tells Redis to enforce the limit exactly, while ~ allows Redis to trim approximately, meaning the stream may temporarily contain more than 10,000 entries.

6 XREADGROUP: Coordinating Consumer Groups

So far, we've seen how to add entries to a stream and read them. But in a real system, we usually have a process responsible for reading these events and doing something with them. This process is called a consumer.

For example, a consumer could read an order.created event and update inventory, send a notification, or process some other downstream work.

Now, what if we have multiple consumers doing the same type of work? We don't want every consumer processing the same event. This is where consumer groups come in.

A consumer group represents a set of consumers that are working on the same type of work. These consumers could be multiple instances of the same service, multiple worker processes, or even multiple workers running inside an application.

Redis distributes new entries among the consumers in the group, so each entry is processed by one consumer within that group.


                         orders
                            │
                     order.created
                            │
                  ┌─────────────────┐
                  │ Consumer Group  │
                  │ order-processors│
                  └────────┬────────┘
                           │
                 ┌─────────┴─────────┐
                 ↓                   ↓
          Order Service         Order Service
           Instance 1            Instance 2
            Consumer 1            Consumer 2
Enter fullscreen mode Exit fullscreen mode

We can create a consumer group using the command XGROUP CREATE

XGROUP CREATE orders order-processors 0
Enter fullscreen mode Exit fullscreen mode

Here, orders is the stream, order-processors is the name of the consumer group, and 0 tells Redis to make the group start from the beginning of the stream.

We can then go on to add consumers;

XREADGROUP GROUP order-processors consumer-1 STREAMS orders >
Enter fullscreen mode Exit fullscreen mode

Here, order-processors is the consumer group, consumer-1 identifies the consumer, and > tells Redis to deliver new entries that have not yet been delivered to another consumer in this group.

Each consumer in the group can then receive different entries from the stream. The entries are still ordered in the stream, but because different consumers process them independently, the processing itself may not happen in the same order so we have to account for that in our application logic when processing events where order matters.

orders
  │
  ├── order.created ──────→ Consumer 1
  │
  ├── order.paid ─────────→ Consumer 2
  │
  ├── order.updated ──────→ Consumer 1
  │
  └── order.cancelled ────→ Consumer 2
Enter fullscreen mode Exit fullscreen mode

Now let's look at what happens when consumer-1 reads from the group.

XREADGROUP GROUP order-processors consumer-1 STREAMS orders >
Enter fullscreen mode Exit fullscreen mode

Redis gives consumer-1 an entry from the stream and records that this consumer has received it. With XREADGROUP, Redis manages the state of the messages delivered to each consumer. This is important because Redis now knows the message has been delivered, but it doesn't yet know whether the consumer has successfully finished processing it.

However, if consumer-1 crashes while processing the work, how do we ensure the message isn't lost?

7. Pending Entries List: Handling Failed Consumers

When a consumer receives an entry using XREADGROUP, Redis adds that entry to the Pending Entries List (PEL) for that consumer group. The entry stays pending until the consumer acknowledges that it has successfully finished processing it.

Consumers acknowledge that they have successfully processed an entry using XACK:

XACK orders order-processors 1706000000000-0
Enter fullscreen mode Exit fullscreen mode

Once acknowledged, Redis removes the entry from the group's PEL. Now we still haven't addressed what happens if the consumer crashes before acknowledging the entry? The entry remains in the PEL because Redis knows it was delivered, but it was never acknowledged as successfully processed and this allows another consumer to take over the pending entry.

Claiming from Pending Entries List

Redis provides XAUTOCLAIM to find pending entries that have been idle for a specified amount of time and transfer them to another consumer.

XAUTOCLAIM orders order-processors consumer-2 60000 0-0
Enter fullscreen mode Exit fullscreen mode

Here, 60000 means the entry must have been pending for at least 60 seconds before another consumer can claim it, we ideally use the idle time because we don't want another consumer to immediately take over an entry that the original consumer is still processing.

Once the new consumer claims the entry, it becomes associated with the new consumer, which can process and acknowledge it normally with XACK.

8. Building a Redis Stream Worker in Go

Now that we understand how Redis Streams, consumer groups, acknowledgements, and pending entries work, let's put everything together and build a simple worker in Go.

We'll build a worker that:

  1. Connects to Redis.
  2. Reads new entries from a consumer group.
  3. Processes the event.
  4. Acknowledges the event after successful processing.
  5. Continues waiting for new events.

First, let's define our Redis connection:

package main

import (
    "context"
    "fmt"
    "log"

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

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

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

    fmt.Println("Connected to Redis")

    _ = rdb
    _ = ctx
}
Enter fullscreen mode Exit fullscreen mode

Before our worker can consume messages, we need to make sure the consumer group exists:

_, err := rdb.XGroupCreateMkStream(
    ctx,
    "orders",
    "order-processors",
    "0",
).Result()

if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
    log.Fatal(err)
}
Enter fullscreen mode Exit fullscreen mode

Here, we're creating the order-processors consumer group for the orders stream. XGroupCreateMkStream also creates the stream if it doesn't already exist.

Now we can continuously read from the stream:

for {
    result, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
        Group:    "order-processors",
        Consumer: "consumer-1",
        Streams:  []string{"orders", ">"},
        Count:    10,
        Block:    0,
    }).Result()

    if err != nil {
        log.Fatal(err)
    }

    for _, stream := range result {
        for _, message := range stream.Messages {
            fmt.Println(message.ID, message.Values)

            // process event

            _, err := rdb.XAck(
                ctx,
                "orders",
                "order-processors",
                message.ID,
            )

            if err != nil {
                log.Printf("failed to acknowledge %s: %v", message.ID, err)
            }
        }
    }
}
Enter fullscreen mode Exit fullscreen mode

Once XReadGroup returns entries, the worker loops through the returned messages and processes each one. We only acknowledge the message after the processing succeeds. XAck tells Redis that the message has been successfully processed. If processing fails or the worker crashes before XAck, the message remains in the Pending Entries List and can later be claimed by another consumer.

I also need to mention that Redis Streams uses a pull-based consumption model. The consumer actively asks Redis for new messages using XREADGROUP, rather than Redis pushing messages directly to the consumer.

9. Redis Streams vs Kafka and RabbitMQ

Redis Streams, Kafka, and RabbitMQ are all message brokers we can use for asynchronous communication, but they each have their own tradeoffs.

Redis Streams

This entire article has been about explaining Redis streams, it is useful when your system already uses Redis and you want stream-based messaging without introducing another infrastructure component.

Redis streams provides us with things like

  • Consumer groups
  • Message acknowledgement
  • Pending Entries List
  • Message replay
  • Blocking reads
  • Configurable retention using MAXLEN and XTRIM

However, Redis Streams has a lighter architecture compared to dedicated streaming platforms like Kafka. Since Streams is a data structure built into Redis, you can use the Redis infrastructure you already have without introducing another dedicated streaming system

Kafka

Kafka is a distributed event streaming platform designed for high-throughput event streaming. Events are stored in topics, which are divided into partitions.

Kafka is also a pull-based system like Redis Streams. Consumers pull records from Kafka, and consumer groups allow multiple consumers to share partitions.

An important difference is how ordering works in Kafka. Kafka guarantees ordering within a partition, but there is no global ordering across multiple partitions.

Kafka is generally used for large-scale event streaming, high throughput, and retaining event history for consumers to replay.

RabbitMQ

RabbitMQ is primarily a message broker built around queues and routing. Producers publish messages to exchanges, and RabbitMQ routes those messages into queues where consumers process them.

RabbitMQ also supports acknowledgements, retries, routing patterns, and multiple exchange types such as direct, topic, fanout, and headers.

Unlike Redis Streams and Kafka, RabbitMQ's consumption model is primarily push-based, meaning once a consumer subscribes to a queue, RabbitMQ delivers messages to that consumer.

So which one should you use?

There isn't a universal choice per se. If you are already using Redis and need a lightweight stream-based communication mechanism, Redis Streams can be a practical choice. I primarily use Redis Streams when I need event-driven communication between services without introducing a separate messaging system.

Kafka becomes useful when event streaming is a major part of the system and you need things like partitioning, high throughput, and long-term event retention.

RabbitMQ becomes useful when you need to distribute tasks across workers, especially when you need reliable delivery and flexible routing.