Настраиваем очереди сообщений на Redis: используем Pub/Sub и Streams
Представьте: ваш интернет-магазин отправляет 10 000 писем в час. Если отправлять синхронно, сервер зависает на минуту, а пользователь ждёт ответа. Асинхронная очередь на Redis решает эту проблему: письма уходят в фоне, а запрос обрабатывается мгновенно. За 5 лет работы мы внедрили Redis очереди в 50+ проектах, сократив нагрузку на серверы до 70%.
Redis Pub/Sub и Streams — два популярных решения для асинхронности. Мы поможем настроить такую инфраструктуру под ключ за 1–2 дня, гарантируем нулевую потерю данных благодаря Streams. Оценим ваш проект бесплатно.
Как выбрать между Pub/Sub и Streams?
Redis предоставляет два механизма для асинхронных сообщений: Pub/Sub — простой fire-and-forget без персистентности, и Streams — персистентная очередь с группами потребителей, похожая на облегчённый Kafka. Выбор зависит от задачи: real-time уведомления (Pub/Sub) или надёжная task queue (Streams).
| Характеристика | Pub/Sub | Streams | Lists (LPUSH/BRPOP) |
|---|---|---|---|
| Персистентность | Нет | Да | Да |
| Consumer groups | Нет | Да | Нет |
| Replay истории | Нет | Да | Нет |
| Сложность | Минимальная | Средняя | Минимальная |
| Производительность на 10K msg/s | 2.1 ms latency | 3.4 ms latency | 1.8 ms latency |
| Применение | Real-time events | Task queue | Simple queue |
Нужна надёжная очередь? Свяжитесь с нами — мы поможем выбрать оптимальный вариант.
Почему Redis Streams лучше Pub/Sub для критичных задач?
Streams в 2–3 раза удобнее Pub/Sub при масштабировании: они поддерживают consumer groups, позволяют подтверждать обработку и перечитывать упавшие сообщения. Pub/Sub — простой вариант для real-time событий, но при росте нагрузки или требованиях к гарантии доставки выбирайте Streams. Экономия на инженерных часах достигает 30%: не нужно писать логику повторной обработки вручную. Кроме того, затраты на инфраструктуру снижаются до 50% за счёт уменьшения числа простоев.
Настройка Redis Pub/Sub
Подходит для real-time уведомлений внутри приложения. Сообщения не сохраняются — если подписчик отключён, сообщение теряется.
// Laravel: публикация через Redis Pub/Sub use Illuminate\Support\Facades\Redis; // Publisher Redis::publish('user-notifications', json_encode([ 'user_id' => $userId, 'type' => 'order.shipped', 'message' => 'Ваш заказ отправлен', ])); // Subscriber (console command) class RedisSubscribeCommand extends Command { protected $signature = 'redis:subscribe'; public function handle(): void { Redis::subscribe(['user-notifications'], function (string $message) { $data = json_decode($message, true); broadcast(new UserNotificationEvent($data)); // → WebSocket }); } } Настройка Redis Streams
Streams — правильный выбор для task queue на Redis. Сообщения хранятся в потоке, consumer groups отслеживают прогресс, pending entries — необработанные сообщения. Гарантия доставки: сообщение удаляется только после XACK.
# Создать поток и добавить сообщение XADD emails * user_id 123 email [email protected] template welcome # Создать consumer group XGROUP CREATE emails email-workers $ MKSTREAM # Читать новые сообщения (воркер 1) XREADGROUP GROUP email-workers worker-1 COUNT 10 BLOCK 5000 STREAMS emails > # Подтвердить обработку XACK emails email-workers <message-id> Как настроить consumer groups в Redis Streams?
Consumer groups позволяют распределять сообщения между воркерами. Каждый воркер получает уникальные сообщения, а pending entries отслеживают необработанные. Это основа отказоустойчивости.
Пример воркера на PHP
use Illuminate\Support\Facades\Redis; class RedisStreamWorker { private string $stream = 'emails'; private string $group = 'email-workers'; private string $consumer; public function __construct() { $this->consumer = gethostname() . ':' . getmypid(); $this->ensureGroup(); } private function ensureGroup(): void { try { Redis::xgroup('CREATE', $this->stream, $this->group, '$', true); } catch (\Throwable) { // Группа уже существует } } public function run(): void { while (true) { // Сначала обработать pending (не подтверждённые с прошлого запуска) $pending = Redis::xreadgroup( $this->group, $this->consumer, [$this->stream => '0'], // '0' = pending messages 10 ); $this->processMessages($pending); // Затем новые сообщения $messages = Redis::xreadgroup( $this->group, $this->consumer, [$this->stream => '>'], // '>' = only new 10, 5000 // блокировка 5 секунд ); $this->processMessages($messages); } } private function processMessages(?array $streams): void { if (!$streams) return; foreach ($streams[$this->stream] ?? [] as [$id, $fields]) { try { $this->handleEmail($fields); Redis::xack($this->stream, $this->group, $id); } catch (\Throwable $e) { Log::error('Stream message failed', ['id' => $id, 'error' => $e->getMessage()]); // Сообщение остаётся в pending — будет перечитано при следующем запуске } } } private function handleEmail(array $fields): void { Mail::to($fields['email'])->send(new TemplateMail($fields['template'], $fields)); } } Пример воркера на Node.js
import Redis from 'ioredis'; const redis = new Redis({ host: 'redis', port: 6379 }); const STREAM = 'emails'; const GROUP = 'email-workers'; const CONSUMER = `worker-${process.pid}`; async function startWorker(): Promise<void> { // Создать группу если не существует try { await redis.xgroup('CREATE', STREAM, GROUP, '$', 'MKSTREAM'); } catch { /* group exists */ } while (true) { const messages = await redis.xreadgroup( 'GROUP', GROUP, CONSUMER, 'COUNT', '10', 'BLOCK', '5000', 'STREAMS', STREAM, '>' ) as [string, [string, string[]][]][] | null; if (!messages) continue; for (const [, entries] of messages) { for (const [id, fields] of entries) { const data = Object.fromEntries( fields.reduce((acc, val, i) => (i % 2 === 0 ? acc.push([val, fields[i+1]]) : acc, acc), [] as [string,string][]) ); try { await sendEmail(data); await redis.xack(STREAM, GROUP, id); } catch (err) { console.error('Email failed:', id, err); } } } } } Сравнение PHP и Node.js для реализации воркеров
| Характеристика | PHP (Laravel) | Node.js (ioredis) |
|---|---|---|
| Параллелизм | Процессы (supervisor) | Event loop |
| Обработка pending | Встроенная (Laravel Horizon) | Ручная |
| Популярность | Широко используется | Высокая производительность |
| Сложность настройки | Средняя | Низкая |
Управление потоком и очистка
# Обрезать поток до 10000 последних сообщений XTRIM emails MAXLEN ~ 10000 # Автоматически при добавлении XADD emails MAXLEN ~ 100000 * user_id 123 template welcome Мониторинг и отладка
Отслеживайте pending entries — необработанные сообщения. Если их количество растёт, воркер не справляется. Используйте XINFO STREAM emails для просмотра состояния. Настройте алерты на длину pending.
Redis Streams — персистентная очередь с группами потребителей. Redis Documentation
Детали мониторинга
- Алерты в Telegram/Slack при превышении порога pending (например, > 1000).
- Логировать ошибки с ID сообщения для ручной повторной обработки.
- Использовать
XCLAIMдля переназначения зависших сообщений другому воркеру.
Типичные ошибки и их решение
- Отсутствие обработки pending: воркер упал, сообщения зависли. Решение — всегда обрабатывать pending при старте.
- Нет гарантии идемпотентности: повторная отправка email. Используйте idempotency key.
- Слишком агрессивный trimming: теряются необработанные сообщения. Используйте ~ (тильда) для приблизительного обрезания.
Что входит в работу и сроки
- Анализ требований и проектирование схемы потоков.
- Реализация воркеров на PHP или Node.js с обработкой ошибок.
- Настройка consumer groups, pending entries и мониторинга.
- Документация по эксплуатации (trimming, алерты).
- Обучение команды (1 час).
- Гарантия на код — 3 месяца.
Базовая реализация Streams воркера (email, уведомления) — 1–2 дня. С мониторингом, алертами и документацией — 2–3 дня. Стоимость рассчитывается индивидуально.
Свяжитесь с нами для бесплатной консультации — оценим ваш проект за 1 день. Закажите настройку Redis очереди уже сегодня! Экономьте ресурсы сервера и время разработчиков.







