Архитектура

Redis как брокер событий: паттерны Pub/Sub и Streams для микросервисов без Kafka

Ruslan Ismailov Опубликовано 12 мин чтения
R

Введение: когда Kafka избыточна и почему Redis стал популярным выбором в 2026 году

Kafka — мощный инструмент, но его развёртывание и эксплуатация требуют значительных ресурсов: отдельный кластер ZooKeeper или KRaft, настройка репликации, мониторинг брокеров, обучение команды. Для стартапов, небольших команд или микросервисных систем с умеренной нагрузкой эти затраты зачастую не оправданы.

В 2026 году Redis всё чаще используется как брокер событий в event-driven архитектурах. Причин несколько: Redis уже есть в большинстве продакшн-стеков как кеш, он хорошо известен разработчикам, прост в операционном обслуживании и предоставляет два зрелых механизма обмена сообщениями — Pub/Sub и Streams. Если ваша команда не имеет выделенного инфраструктурного опыта, а объёмы сообщений исчисляются тысячами, а не миллионами в секунду, Redis как альтернатива Kafka заслуживает серьёзного рассмотрения.

Redis Pub/Sub: механика работы и гарантии доставки

Как работает Pub/Sub

Redis Pub/Sub — это классическая модель «публикация-подписка» в памяти. Издатель отправляет сообщение в канал командой PUBLISH, все активные подписчики этого канала мгновенно получают его через постоянное соединение. Никакого буферизирования: если подписчик не подключён в момент публикации, сообщение теряется безвозвратно.

Redis поддерживает как точные каналы, так и паттерн-подписку через PSUBSCRIBE, что позволяет подписываться на группы каналов по маске, например orders.*.

Гарантии доставки — их отсутствие

Pub/Sub в Redis работает по принципу fire-and-forget: нет персистентности, нет подтверждений, нет повторной доставки. Это не баг, а особенность дизайна. Если консьюмер упал, перезапустился или просто был медленнее продюсера — сообщения потеряны. Такая семантика называется «at most once».

Когда Pub/Sub уместен

  • Инвалидация кеша: уведомления об изменении данных, где потеря одного события не критична
  • Real-time нотификации в UI: онлайн-статус пользователей, чаты
  • Широковещательные события без гарантий: системные метрики, heartbeat-сообщения
  • Прототипирование event-driven потоков до внедрения более сложного решения

Redis Streams: архитектура, consumer groups и ACK-механизм

Что такое Redis Streams

Redis Streams, представленные в Redis 5.0, — это персистентная, упорядоченная структура данных, напоминающая append-only лог. Каждая запись имеет уникальный ID вида 1704067200000-0 (миллисекунды + порядковый номер). Стримы хранятся в памяти и могут реплицироваться и персистироваться через стандартные механизмы RDB/AOF.

Consumer Groups

Ключевое отличие от Pub/Sub — поддержка consumer groups. Группа консьюмеров обрабатывает поток совместно: каждое сообщение доставляется ровно одному члену группы. Это позволяет горизонтально масштабировать обработку: несколько воркеров читают из одного стрима, не дублируя работу. Разные группы получают все сообщения независимо — аналог топиков Kafka с несколькими группами потребителей.

ACK-механизм и pending entries

После получения сообщения консьюмер должен явно подтвердить обработку командой XACK. До подтверждения сообщение остаётся в Pending Entries List (PEL). Если консьюмер падает, необработанные сообщения можно переназначить другому воркеру командой XCLAIM. Это обеспечивает семантику «at least once».

Персистентность и ретенция

Стримы поддерживают ограничение длины через MAXLEN — можно хранить последние N сообщений или обрезать по времени. При использовании AOF или RDB данные переживают перезапуск Redis. Однако важно понимать: Redis — это in-memory база данных, и при полной потере данных (без репликации и персистентности) стрим будет утрачен.

Сравнение: Redis Pub/Sub vs Redis Streams vs Kafka

Ниже — ключевые параметры трёх решений, которые важны при выборе брокера событий для микросервисной архитектуры.

  • Персистентность: Pub/Sub — нет; Streams — да (RDB/AOF); Kafka — да (дисковый лог)
  • Гарантии доставки: Pub/Sub — at most once; Streams — at least once; Kafka — at least once / exactly once
  • Consumer groups: Pub/Sub — нет; Streams — да; Kafka — да
  • Повторная обработка: Pub/Sub — невозможна; Streams — да, через XCLAIM; Kafka — да, сброс offset
  • Порядок сообщений: Pub/Sub — не гарантирован; Streams — гарантирован в рамках стрима; Kafka — гарантирован в рамках партиции
  • Операционная сложность: Pub/Sub — минимальная; Streams — низкая; Kafka — высокая
  • Пропускная способность: Pub/Sub — высокая; Streams — высокая (сотни тысяч/сек); Kafka — очень высокая (миллионы/сек)
  • Долгосрочное хранение: Pub/Sub — нет; Streams — ограниченное (RAM); Kafka — практически неограниченное
  • Инфраструктурные зависимости: Pub/Sub — только Redis; Streams — только Redis; Kafka — Kafka + ZooKeeper/KRaft

Паттерн Event-Driven с Redis Streams: реализация на Go

Продюсер

Пример продюсера на Go с использованием библиотеки go-redis/v9. Продюсер публикует событие об оформлении заказа в стрим orders.

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()

    // Публикуем событие в стрим
    id, err := rdb.XAdd(ctx, &redis.XAddArgs{
        Stream: "orders",
        MaxLen: 10000,    // ограничение длины стрима
        Approx: true,     // приблизительная обрезка для производительности
        Values: map[string]interface{}{
            "order_id":   "ORD-12345",
            "user_id":    "USR-67890",
            "amount":     "1500.00",
            "currency":   "RUB",
            "created_at": time.Now().Unix(),
        },
    }).Result()

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

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

Консьюмер с consumer group

Консьюмер создаёт группу (если её нет), читает сообщения и подтверждает обработку через 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"))

    // Создаём группу, если не существует
    err := rdb.XGroupCreateMkStream(ctx, streamName, groupName, "$").Err()
    if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
        log.Fatalf("XGroupCreate error: %v", err)
    }

    // Обработка сигналов для 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:
        }

        // Читаем до 10 сообщений с таймаутом 2 секунды
        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 // таймаут, нет новых сообщений
            }
            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)
                    // Сообщение останется в PEL для повторной обработки
                    continue
                }

                // Подтверждаем успешную обработку
                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"],
    )
    // Бизнес-логика обработки заказа
    return nil
}

Обработка ошибок и Dead Letter: что делать с необработанными сообщениями

Одна из типичных проблем при работе с Redis Streams — «зависшие» сообщения в PEL. Сообщение может оказаться там, если консьюмер упал до вызова XACK или если обработка постоянно завершается ошибкой.

Стратегия обработки PEL

Рекомендуемый паттерн: периодически запускать задачу, которая через XPENDING проверяет сообщения, «зависшие» дольше порогового времени (например, 5 минут), и перемещает их на другой консьюмер через XCLAIM. После N неудачных попыток сообщение переносится в отдельный стрим-«мусорную корзину» — 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 {
            // Перемещаем в dead letter stream
            moveToDeadLetter(ctx, rdb, p.ID)
            rdb.XAck(ctx, streamName, groupName, p.ID)
            continue
        }

        // Переназначаем текущему воркеру
        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)
}

Масштабирование консьюмеров в Kubernetes

Деплой воркеров

Каждый Pod-воркер в Kubernetes использует переменную окружения POD_NAME как уникальный идентификатор консьюмера в группе. Redis сам балансирует нагрузку между членами группы — никакой дополнительной координации не требуется.

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

Автомасштабирование по длине стрима

KEDA (Kubernetes Event-Driven Autoscaling) поддерживает Redis Streams из коробки. Скейлер читает длину PEL или количество необработанных сообщений и масштабирует Deployment автоматически.

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"  # масштабируем, если PEL > 50

Мониторинг Redis как брокера: ключевые метрики и инструменты

Ключевые метрики

  • Длина стрима (XLEN orders) — растущая длина сигнализирует о том, что консьюмеры не справляются
  • Размер PEL (XPENDING orders order-processor - + 1) — количество необработанных сообщений
  • Lag группы — разница между последним ID в стриме и последним подтверждённым ID группы
  • used_memory — Redis хранит всё в RAM; переполнение памяти критично
  • connected_clients — число активных подключений консьюмеров
  • evicted_keys — если Redis вытесняет ключи из-за нехватки памяти, данные могут теряться

Инструменты

Для мониторинга Redis как брокера событий используйте связку Redis Exporter + Prometheus + Grafana. Redis Exporter (oliver006/redis_exporter) экспортирует метрики redis_stream_length, redis_stream_groups и redis_stream_pending_entries_count. Настройте алерты на: длина стрима > 10 000, PEL > 500, lag группы > 60 секунд, used_memory > 80% от maxmemory.

Ограничения Redis как брокера: когда всё-таки нужна Kafka

Честность важна: Redis как брокер событий имеет принципиальные ограничения, о которых нужно знать заранее.

  • Объём хранения ограничен RAM. Kafka хранит данные на диске и может удерживать сообщения неделями. Redis при нехватке памяти начинает вытеснять данные — и сообщения теряются, если не настроен MAXLEN и не контролируется память.
  • Нет нативного партиционирования. Один стрим — одна последовательность. Для масштабирования нужно вручную создавать несколько стримов и распределять между ними нагрузку.
  • Репликация не синхронная. В Redis Sentinel и Cluster репликация асинхронная. При failover возможна потеря нескольких последних сообщений.
  • Нет exactly-once семантики. Redis Streams обеспечивают at least once. Для exactly once нужна идемпотентность на уровне приложения.
  • Нет встроенного Schema Registry. Совместимость форматов сообщений — забота разработчика.

Выбирайте Kafka, если: вам нужно хранить события месяцами, объём превышает несколько миллионов сообщений в секунду, требуется строгая exactly-once семантика, или у вас есть выделенная инфраструктурная команда для его обслуживания.

Рассмотрите RabbitMQ, если нужна более богатая маршрутизация сообщений (exchanges, routing keys, dead letter queues из коробки) при умеренных объёмах.

Реальный кейс: замена Redis Pub/Sub на Redis Streams в продакшне

Контекст

Команда e-commerce платформы использовала Redis Pub/Sub для уведомлений об изменении статуса заказов. Микросервис уведомлений подписывался на канал order.status.changed и отправлял push-уведомления и email. Система работала, но периодически возникали инциденты: при перезапуске Pod-а консьюмера в Kubernetes часть уведомлений терялась — пользователи не получали письма о доставке.

Проблема

Pub/Sub не буферизирует сообщения. Пока Pod перезапускался (30–60 секунд), все события терялись. Rolling update с нулевым даунтаймом не помогал: в момент смены Pod-ов всегда был краткий разрыв подписки.

Решение

Команда перешла на Redis Streams с consumer group notification-service. Продюсер (сервис управления заказами) стал писать в стрим order-events вместо публикации в канал. Консьюмер читает через XREADGROUP и подтверждает XACK только после успешной отправки уведомления. PEL-воркер раз в минуту проверяет зависшие сообщения и повторяет обработку.

Результат

За три месяца после миграции команда не зафиксировала ни одного инцидента с потерей уведомлений. Время миграции составило два дня, включая тестирование. Дополнительных инфраструктурных зависимостей не появилось — Redis уже был в стеке. Операционные затраты на поддержку остались прежними.

Итог и рекомендации по выбору

Redis как брокер событий — это прагматичный выбор для команд, которые хотят получить event-driven архитектуру без операционной сложности Kafka.

  • Используйте Redis Pub/Sub для broadcast-уведомлений, инвалидации кеша и real-time событий, где потеря одного-двух сообщений не критична.
  • Используйте Redis Streams для надёжной асинхронной обработки задач, межсервисного обмена событиями и любого сценария, где важна гарантия доставки at least once.
  • Добавьте KEDA для автомасштабирования консьюмеров в Kubernetes по длине стрима или PEL.
  • Всегда настраивайте MAXLEN на стримах и мониторьте used_memory — это защитит от потери данных при нехватке RAM.
  • Реализуйте dead letter stream и логику XCLAIM для корректной обработки сбоев.
  • Если ваша нагрузка выросла до миллионов событий в секунду, требуется долгосрочное хранение или exactly-once семантика — переходите на Kafka. Redis сам даст вам понять, когда это произойдёт.

Redis Streams — это не «бедная версия Kafka». Это другой инструмент с другими trade-off-ами, который идеально подходит для большинства микросервисных систем с умеренной нагрузкой.

Технологии

Теги

Руслан Исмаилов

Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →