Архитектура

Redis Streams как замена очередям сообщений: практическое руководство с примерами

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

Введение: почему Redis Streams — это не просто Pub/Sub

Большинство разработчиков знают Redis как высокопроизводительный кэш. Многие используют Redis Pub/Sub для простых уведомлений. Но Redis Streams — это отдельная, значительно более мощная структура данных, появившаяся в Redis 5.0, которая кардинально меняет подход к организации очередей сообщений в микросервисной архитектуре.

Чем Redis Streams принципиально отличается от Pub/Sub? В Pub/Sub сообщения не персистируются: если подписчик недоступен в момент публикации, он пропустит сообщение навсегда. Redis Streams, напротив, хранят все сообщения в упорядоченном журнале с уникальными ID, поддерживают consumer groups с подтверждением доставки (ACK) и позволяют переигрывать историю сообщений — почти как Apache Kafka, но без его инфраструктурной сложности.

Эта статья предназначена для backend-разработчиков, которые уже работают с Redis как кэшем, и для архитекторов микросервисов, которые ищут прагматичную альтернативу RabbitMQ или Kafka для задач средней сложности. Мы разберём всё: от базовых концепций до production-ready кода на Go и интеграции в Laravel.

Основные концепции Redis Streams

Stream — журнал сообщений

Stream в Redis — это append-only log: структура данных, в которую можно только добавлять новые записи (entries). Каждая запись имеет уникальный ID формата millisecondsTime-sequenceNumber (например, 1700000000000-0) и набор произвольных полей ключ-значение.

# Добавить сообщение в stream
XADD orders * user_id 42 product_id 101 action purchase

# Прочитать последние 10 сообщений
XRANGE orders - + COUNT 10

# Получить длину stream
XLEN orders

Символ * в команде XADD означает «сгенерируй ID автоматически». Вы также можете передать явный ID, что полезно при репликации данных из внешних систем.

Consumer Groups — параллельная обработка

Consumer Group — это механизм, позволяющий нескольким консьюмерам совместно обрабатывать один stream, при этом каждое сообщение доставляется только одному консьюмеру в группе. Это ключевое отличие от обычного XREAD, где каждый читатель получает все сообщения.

# Создать consumer group
XGROUP CREATE orders processing-group $ MKSTREAM

# Прочитать новые сообщения как консьюмер worker-1
XREADGROUP GROUP processing-group worker-1 COUNT 5 BLOCK 2000 STREAMS orders >

# Подтвердить обработку сообщений (ACK)
XACK orders processing-group 1700000000000-0 1700000000001-0

Символ > означает «дай мне только необработанные сообщения». Символ $ при создании группы означает «начинай с новых сообщений»; 0 — с самого начала stream.

Pending Entries List (PEL)

Когда консьюмер забирает сообщение через XREADGROUP, оно попадает в Pending Entries List — список сообщений, доставленных, но ещё не подтверждённых. Это критически важно для гарантии доставки: если воркер упал, не успев отправить ACK, сообщение останется в PEL и может быть передано другому воркеру.

# Посмотреть pending сообщения
XPENDING orders processing-group - + 10

# Детальная информация по конкретным pending записям
XPENDING orders processing-group - + 10 worker-1

Сравнение Redis Streams с RabbitMQ и Kafka

Прежде чем писать код, важно понять, когда Redis Streams — правильный выбор, а когда лучше взять специализированный брокер.

КритерийRedis StreamsRabbitMQApache Kafka
ПерсистентностьRDB/AOF, опциональноDisk (durable queues)Disk (partition log)
ThroughputВысокий (сотни тыс/с)СреднийОчень высокий (млн/с)
ЗадержкаОчень низкая (<1 мс)НизкаяСредняя (batch)
Replay сообщенийДа (по offset)НетДа
RoutingПростой (по имени stream)Сложный (exchanges, bindings)По топикам/партициям
Операционная сложностьНизкаяСредняяВысокая
Гарантии доставкиAt-least-onceAt-least-once / exactly-onceAt-least-once / exactly-once

Используйте Redis Streams, когда:

  • У вас уже есть Redis в инфраструктуре, и вы хотите избежать нового компонента
  • Нагрузка — тысячи, а не миллионы сообщений в секунду
  • Нужна простая событийная архитектура для микросервисов
  • Важна низкая задержка
  • Не нужна сложная маршрутизация

Оставайтесь на Kafka, если:

  • Объёмы данных — сотни гигабайт в день
  • Нужна многоуровневая репликация и гарантии exactly-once
  • Требуется долгосрочное хранение событий (event sourcing в масштабе)

Оставайтесь на RabbitMQ, если:

  • Нужна сложная маршрутизация через exchanges и bindings
  • Важна поддержка AMQP-протокола
  • Команда хорошо знает RabbitMQ и нет смысла в миграции

Практика: реализация producer/consumer на Go

Рассмотрим реальный пример: система обработки заказов в интернет-магазине. Producer публикует события создания заказов, несколько consumer воркеров обрабатывают их параллельно.

Установка зависимостей

go mod init orders-processor
go get github.com/redis/go-redis/v9

Producer: публикация событий

package main

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

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

func main() {
    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    ctx := context.Background()

    // Создаём consumer group при старте (игнорируем ошибку если уже существует)
    err := rdb.XGroupCreateMkStream(ctx, "orders", "processing-group", "$").Err()
    if err != nil && err.Error() != "BUSYGROUP Consumer Group name already exists" {
        log.Fatalf("Failed to create consumer group: %v", err)
    }

    // Публикуем события заказов
    for i := 1; i <= 100; i++ {
        msgID, err := rdb.XAdd(ctx, &redis.XAddArgs{
            Stream: "orders",
            Values: map[string]interface{}{
                "order_id":   fmt.Sprintf("ORD-%04d", i),
                "user_id":    i * 10,
                "amount":     float64(i) * 99.9,
                "created_at": time.Now().Unix(),
            },
        }).Result()
        if err != nil {
            log.Printf("Failed to publish order: %v", err)
            continue
        }
        fmt.Printf("Published order ORD-%04d with ID: %s\n", i, msgID)
        time.Sleep(10 * time.Millisecond)
    }
}

Consumer: обработка сообщений с ACK

package main

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

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

const (
    streamName    = "orders"
    groupName     = "processing-group"
    blockDuration = 2 * time.Second
    maxRetries    = 3
)

func main() {
    consumerName := os.Getenv("CONSUMER_NAME")
    if consumerName == "" {
        consumerName = "worker-1"
    }

    rdb := redis.NewClient(&redis.Options{
        Addr: "localhost:6379",
    })
    ctx, cancel := context.WithCancel(context.Background())

    // Graceful shutdown
    sigCh := make(chan os.Signal, 1)
    signal.Notify(sigCh, syscall.SIGTERM, syscall.SIGINT)
    go func() {
        <-sigCh
        fmt.Println("Shutting down...")
        cancel()
    }()

    // Сначала обрабатываем pending сообщения (после перезапуска)
    processPendingMessages(ctx, rdb, consumerName)

    // Основной цикл обработки новых сообщений
    for {
        select {
        case <-ctx.Done():
            return
        default:
            messages, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
                Group:    groupName,
                Consumer: consumerName,
                Streams:  []string{streamName, ">"},
                Count:    10,
                Block:    blockDuration,
            }).Result()
            if err != nil {
                if err == redis.Nil {
                    continue // timeout, нет новых сообщений
                }
                log.Printf("Error reading from stream: %v", err)
                time.Sleep(time.Second)
                continue
            }

            for _, stream := range messages {
                for _, msg := range stream.Messages {
                    if err := processOrder(msg); err != nil {
                        log.Printf("Failed to process message %s: %v", msg.ID, err)
                        continue
                    }
                    // Подтверждаем успешную обработку
                    if err := rdb.XAck(ctx, streamName, groupName, msg.ID).Err(); err != nil {
                        log.Printf("Failed to ACK message %s: %v", msg.ID, err)
                    }
                }
            }
        }
    }
}

func processOrder(msg redis.XMessage) error {
    orderID := msg.Values["order_id"]
    amount := msg.Values["amount"]
    fmt.Printf("Processing order %s, amount: %s\n", orderID, amount)
    // Здесь ваша бизнес-логика: сохранение в БД, отправка email и т.д.
    time.Sleep(50 * time.Millisecond) // эмуляция работы
    return nil
}

func processPendingMessages(ctx context.Context, rdb *redis.Client, consumerName string) {
    pending, err := rdb.XPendingExt(ctx, &redis.XPendingExtArgs{
        Stream:   streamName,
        Group:    groupName,
        Start:    "-",
        Stop:     "+",
        Count:    100,
        Consumer: consumerName,
    }).Result()
    if err != nil {
        return
    }
    for _, p := range pending {
        if p.RetryCount >= maxRetries {
            log.Printf("Message %s exceeded retry limit, moving to DLQ", p.ID)
            // Логика Dead Letter Queue
            continue
        }
        // Reclaim и обработать
        msgs, err := rdb.XClaim(ctx, &redis.XClaimArgs{
            Stream:   streamName,
            Group:    groupName,
            Consumer: consumerName,
            MinIdle:  30 * time.Second,
            Messages: []string{p.ID},
        }).Result()
        if err != nil {
            continue
        }
        for _, msg := range msgs {
            if err := processOrder(msg); err == nil {
                rdb.XAck(ctx, streamName, groupName, msg.ID)
            }
        }
    }
}

Обработка ошибок и гарантии доставки

ACK и гарантия at-least-once

Redis Streams предоставляет гарантию доставки at-least-once: сообщение считается обработанным только после явного XACK. Если воркер упал до ACK, сообщение остаётся в PEL и может быть перезаявлено (XCLAIM) другим воркером.

XPENDING: инспекция необработанных сообщений

# Сводная информация по pending
XPENDING orders processing-group - + 100

# Результат:
# 1) "1700000001234-0"
# 2) "worker-1"
# 3) (integer) 85000     <-- время с последней доставки в мс
# 4) (integer) 2         <-- количество попыток доставки

XCLAIM: перехват зависших сообщений

XCLAIM позволяет одному воркеру перехватить сообщение, которое другой воркер держит слишком долго. Это критично для обеспечения надёжности при сбоях воркеров.

# Перехватить сообщения, idle более 60 секунд
XCLAIM orders processing-group worker-2 60000 1700000001234-0

# Автоматически перехватить batch зависших сообщений (Redis 6.2+)
XAUTOCLAIM orders processing-group worker-2 60000 0-0 COUNT 10

Dead Letter Queue (DLQ)

Для сообщений, которые не удалось обработать после N попыток, рекомендуется реализовать Dead Letter Queue — отдельный stream для «отравленных» сообщений:

# Перенести проблемное сообщение в DLQ
XADD orders-dlq * original_id 1700000001234-0 order_id ORD-0042 reason "processing_failed" attempts 3
XACK orders processing-group 1700000001234-0
XDEL orders 1700000001234-0

Ограничение размера stream

Redis хранит все сообщения в памяти (при работе без офлоада на диск). Чтобы избежать переполнения, используйте параметр MAXLEN:

# Ограничить stream до 100 000 сообщений (приближённое)
XADD orders MAXLEN ~ 100000 * order_id ORD-0001 ...

# Или периодически обрезать вручную
XTRIM orders MAXLEN ~ 100000

Интеграция Redis Streams в Laravel-приложение

Laravel имеет встроенный драйвер Redis для очередей, но он работает через RPUSH/LPOP (List), а не через Streams. Покажем, как создать кастомный драйвер Streams для Laravel Queue.

Создание драйвера

<?php
// app/Queue/RedisStreamConnector.php
namespace App\Queue;

use Illuminate\Queue\Connectors\ConnectorInterface;

class RedisStreamConnector implements ConnectorInterface
{
    public function connect(array $config): RedisStreamQueue
    {
        return new RedisStreamQueue(
            app('redis')->connection($config['connection'] ?? 'default'),
            $config['queue'] ?? 'default',
            $config['group'] ?? 'laravel-workers',
            $config['consumer'] ?? gethostname(),
        );
    }
}
<?php
// app/Queue/RedisStreamQueue.php
namespace App\Queue;

use Illuminate\Contracts\Queue\Queue;
use Illuminate\Queue\Queue as BaseQueue;
use Illuminate\Redis\Connections\Connection;

class RedisStreamQueue extends BaseQueue implements Queue
{
    public function __construct(
        protected Connection $redis,
        protected string $default,
        protected string $group,
        protected string $consumer,
    ) {}

    public function push($job, $data = '', $queue = null): mixed
    {
        $queue = $this->getQueue($queue);
        $payload = $this->createPayload($job, $queue, $data);

        return $this->redis->command('xadd', [
            $queue, '*',
            'payload', $payload,
            'attempts', 0,
        ]);
    }

    public function pop($queue = null): ?\Illuminate\Contracts\Queue\Job
    {
        $queue = $this->getQueue($queue);
        $this->ensureGroupExists($queue);

        $results = $this->redis->command('xreadgroup', [
            'GROUP', $this->group, $this->consumer,
            'COUNT', 1,
            'BLOCK', 2000,
            'STREAMS', $queue, '>',
        ]);

        if (empty($results[$queue])) {
            return null;
        }

        [$messageId, $values] = $results[$queue][0];
        $payload = json_decode($values['payload'], true);

        return new RedisStreamJob(
            $this->container, $this, $this->redis,
            $queue, $this->group, $messageId, $payload,
        );
    }

    public function ack(string $queue, string $messageId): void
    {
        $this->redis->command('xack', [$queue, $this->group, $messageId]);
    }

    protected function ensureGroupExists(string $queue): void
    {
        try {
            $this->redis->command('xgroup', ['CREATE', $queue, $this->group, '$', 'MKSTREAM']);
        } catch (\Exception $e) {
            // Группа уже существует — это нормально
        }
    }

    protected function getQueue(?string $queue): string
    {
        return 'stream:' . ($queue ?? $this->default);
    }

    // Остальные обязательные методы интерфейса Queue...
    public function size($queue = null): int { return 0; }
    public function later($delay, $job, $data = '', $queue = null): mixed { return null; }
    public function bulk($jobs, $data = '', $queue = null): void {}
}

Регистрация драйвера в AppServiceProvider

<?php
// app/Providers/AppServiceProvider.php
public function boot(): void
{
    Queue::extend('redis-stream', function () {
        return new \App\Queue\RedisStreamConnector();
    });
}

Конфигурация в config/queue.php

'connections' => [
    'redis-stream' => [
        'driver'     => 'redis-stream',
        'connection' => 'default',
        'queue'      => 'default',
        'group'      => 'laravel-workers',
        'consumer'   => env('QUEUE_CONSUMER_NAME', gethostname()),
        'retry_after' => 90,
    ],
],

После этого все стандартные Laravel-джобы (dispatch, Queue::push) будут работать через Redis Streams, а вы получаете все преимущества: ACK, PEL, replay и consumer groups.

Масштабирование: несколько Consumer Groups и партиционирование

Несколько Consumer Groups для разных целей

Одна из мощнейших возможностей Redis Streams: один stream может читаться несколькими независимыми consumer groups. Это позволяет реализовать паттерн fan-out без дублирования данных:

# Одни и те же события заказов читают три разных сервиса
XGROUP CREATE orders notification-service $ MKSTREAM
XGROUP CREATE orders analytics-service $ MKSTREAM  
XGROUP CREATE orders inventory-service $ MKSTREAM

Каждая группа получает все сообщения и обрабатывает их независимо. Это фундаментально отличается от round-robin внутри одной группы, где каждое сообщение получает только один воркер.

Партиционирование через несколько streams

Redis — single-threaded по умолчанию, поэтому для горизонтального масштабирования используют несколько именованных streams (аналог партиций в Kafka):

# Producer определяет партицию по хешу user_id
func getPartition(userID int, numPartitions int) string {
    return fmt.Sprintf("orders:partition:%d", userID % numPartitions)
}

// Публикация в партицию
partition := getPartition(userID, 8) // 8 партиций
rdb.XAdd(ctx, &redis.XAddArgs{
    Stream: partition,
    Values: orderData,
})

При использовании Redis Cluster партиции автоматически распределяются по разным нодам, обеспечивая линейное масштабирование throughput.

Автоматическое масштабирование воркеров

В Kubernetes можно настроить HPA (Horizontal Pod Autoscaler) на основе длины PEL или lag consumer group. Метрики экспортируются через redis_exporter в Prometheus:

# Пример HPA на основе custom metrics
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: order-worker-hpa
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: order-worker
  minReplicas: 2
  maxReplicas: 20
  metrics:
  - type: External
    external:
      metric:
        name: redis_stream_pending_entries
        selector:
          matchLabels:
            stream: orders
      target:
        type: AverageValue
        averageValue: "100"

Мониторинг Redis Streams: XINFO, метрики, алерты

XINFO: встроенная диагностика

# Общая информация о stream
XINFO STREAM orders

# Информация о consumer groups
XINFO GROUPS orders

# Информация о консьюмерах в группе
XINFO CONSUMERS orders processing-group

# Полная информация (Redis 7.0+)
XINFO STREAM orders FULL COUNT 10

Ключевые метрики для мониторинга:

  • pending-messages — количество сообщений в PEL (должно быть близко к 0 в норме)
  • lag — разница между последним ID в stream и последним прочитанным ID группы
  • consumers count — количество активных воркеров
  • idle time — время бездействия консьюмера

Мониторинг через redis_exporter + Prometheus

Используйте redis_exporter для экспорта метрик в Prometheus. Добавьте в конфигурацию мониторинг stream-метрик:

# prometheus.yml
scrape_configs:
  - job_name: 'redis'
    static_configs:
      - targets: ['redis-exporter:9121']
    params:
      stream-groups:
        - orders
        - payments

Алерты в Alertmanager

# Алерт на высокий lag
- alert: RedisStreamHighLag
  expr: redis_stream_group_lag{stream="orders"} > 1000
  for: 5m
  labels:
    severity: warning
  annotations:
    summary: "Redis Stream lag is too high"
    description: "Consumer group lag: {{ $value }} messages"

# Алерт на зависшие pending сообщения
- alert: RedisStreamStalePending  
  expr: redis_stream_group_pending{stream="orders"} > 100
  for: 10m
  labels:
    severity: critical

Grafana Dashboard

Рекомендуемые панели для дашборда Redis Streams:

  1. Stream length (XLEN) в динамике
  2. Consumer group lag по каждой группе
  3. Pending entries count
  4. ACK rate (сообщений/сек)
  5. Oldest pending message age
  6. Consumer count per group

Заключение

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

Главные выводы из статьи:

  • Redis Streams отличается от Pub/Sub персистентностью и гарантиями доставки через ACK/PEL
  • Consumer Groups позволяют горизонтально масштабировать обработку сообщений
  • XCLAIM и XAUTOCLAIM решают проблему зависших сообщений при сбоях воркеров
  • Интеграция в Go занимает несколько часов, в Laravel — можно реализовать кастомный драйвер очередей
  • Мониторинг lag и pending entries — ключ к надёжной работе в production

Redis Streams — не замена Kafka для petabyte-scale систем, но отличная альтернатива RabbitMQ для большинства микросервисных задач. Начните с одного stream, одной consumer group и двух воркеров — и вы удивитесь, насколько далеко это вас заведёт.

Технологии

Теги

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

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