Redis Streams как замена очередям сообщений: практическое руководство с примерами
Введение: почему 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 Streams | RabbitMQ | Apache Kafka |
|---|---|---|---|
| Персистентность | RDB/AOF, опционально | Disk (durable queues) | Disk (partition log) |
| Throughput | Высокий (сотни тыс/с) | Средний | Очень высокий (млн/с) |
| Задержка | Очень низкая (<1 мс) | Низкая | Средняя (batch) |
| Replay сообщений | Да (по offset) | Нет | Да |
| Routing | Простой (по имени stream) | Сложный (exchanges, bindings) | По топикам/партициям |
| Операционная сложность | Низкая | Средняя | Высокая |
| Гарантии доставки | At-least-once | At-least-once / exactly-once | At-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/v9Producer: публикация событий
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 10Dead 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: criticalGrafana Dashboard
Рекомендуемые панели для дашборда Redis Streams:
- Stream length (XLEN) в динамике
- Consumer group lag по каждой группе
- Pending entries count
- ACK rate (сообщений/сек)
- Oldest pending message age
- 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. Подробнее обо мне →