В Valkey есть два способа передавать сообщения между процессами, и они дают противоположные обещания. Pub/sub работает по принципу «отправил и забыл»: сообщение уходит тем, кто подписан в этот момент, и затем исчезает - без хранения, без подтверждения и без повторного проигрывания. Streams - это журнал только для добавления: сообщения хранятся, у каждого есть id, потребители могут читать с любого места, а группы потребителей отслеживают, какие сообщения обработал каждый исполнитель и какие он ещё не подтвердил.
Правило, которое из этого следует, короткое. Используйте pub/sub, когда пропущенное сообщение неважно, потому что следующее его заменит, - живые уведомления, подсказки об инвалидации кэша, присутствие в сети. Используйте streams, когда каждое сообщение должно быть обработано хотя бы один раз, - заказы, события, обновляющие другие системы, всё, что вам было бы обидно потерять. Большинство ошибок и с тем, и с другим возникает от использования pub/sub там, где нужен был stream. В этой статье оба механизма разобраны как следует: команды, гарантии, группы потребителей, список ожидающих записей, обрезка и сценарии отказов.
Разница в одной таблице#
| Pub/sub | Streams | |
|---|---|---|
| Хранится | Нет | Да, до обрезки |
| Доставка | Не более одного раза, текущим подписчикам | Хотя бы один раз с группами потребителей |
| Поздний подписчик видит старые сообщения | Нет | Да, с любого id |
| Подтверждение | Нет | XACK для каждого сообщения |
| Балансировка между исполнителями | Нет, каждый подписчик получает каждое сообщение | Да, внутри группы потребителей |
| Переживает перезапуск Valkey | Переживать нечего | Да, если включено сохранение |
| Память | Только буферы для медленных подписчиков | Растёт с числом хранимых записей |
| Типичное применение | Живые обновления, инвалидация, рассылка в чатах | Журналы событий, распределение работы, аудит |
Есть и третий вариант, о котором забывают: обычный список с LPUSH и BRPOP, на котором построены простые очереди задач. Он доставляет каждое сообщение одному потребителю и хранит его, пока его не заберут, но как только потребитель снял сообщение, падение его теряет. Streams с группами потребителей исправляют ровно это, ради чего они и существуют. Полноценные библиотеки очередей задач на этих примитивах описаны в статье очереди задач на Valkey.
Pub/sub: команды и поведение#
Издатель отправляет сообщение в канал; каждое подключение, подписанное в этот момент на канал, его получает.
# connection ASUBSCRIBE orders:updatesPSUBSCRIBE chat:room:*# connection BPUBLISH orders:updates '{"id":512,"status":"shipped"}'(integer) 1PUBLISH возвращает число подписчиков, получивших сообщение. Ответ 0 означает, что сообщение ушло в никуда - никто не слушал, и оно пропало. PSUBSCRIBE подписывает на glob-шаблон; это удобно, но каждое опубликованное сообщение сверяется с каждой подпиской по шаблону, так что тысячи шаблонов стоят CPU при каждой публикации.
Что удивляет людей:
- Подписанное подключение - особое. В протоколе RESP2 подключение в режиме подписки может выполнять только команды, связанные с подпиской, и
PING. Поэтому клиентские библиотеки используют для подписки отдельное подключение. Не пытайтесь делить с ней ваш основной клиент. - Отключился - значит пропустил. Подписчик, переподключившийся после сетевого сбоя, deploy или паузы сборщика мусора, не получает ничего из опубликованного, пока его не было. Никакого догоняния нет.
- Медленных подписчиков отключают. Сообщения, ждущие подписчика, буферизуются в памяти Valkey.
client-output-buffer-limitдля клиентов pub/sub по умолчанию равен32mb 8mb 60: подписчик, чей буфер превысил 32 МБ или держится выше 8 МБ в течение 60 секунд, отключается. Это защищает сервер от одного застрявшего потребителя и означает, что медленный подписчик теряет сообщения. - Pub/sub не хранится, поэтому сохранение данных на него никак не влияет. AOF и снимки сохраняют ключи, а каналы - это не ключи.
Шардированный pub/sub (SSUBSCRIBE, SPUBLISH) существует для кластерных развёртываний, где он удерживает сообщения в пределах одного шарда. На одиночном сервере он ведёт себя как обычный pub/sub.
Для чего хорош pub/sub#
С такими свойствами pub/sub подходит для работы, где потеря сообщения безвредна:
- Рассылка в реальном времени в браузеры. Несколько веб-процессов держат каждый свои подключения WebSocket; событие, опубликованное один раз, доходит до каждого процесса, а тот пересылает его своим подключённым пользователям. Так работают и Redis-адаптер Socket.IO, и Redis backplane в ASP.NET Core SignalR - эти настройки описаны в статьях приложения реального времени на SignalR и WebSockets за обратным прокси. Если пользователь пропустит сигнал «новое сообщение», он увидит сообщение при следующем обновлении.
- Подсказки об инвалидации кэша. Процессы с кэшем в памяти подписываются на
cache:invalidate; запись публикует изменившийся ключ. Сочетайте это с TTL в локальном кэше, чтобы пропущенная подсказка означала лишь ненадолго устаревшие данные. - Присутствие в сети и индикаторы набора текста. По определению важны только тем, кто слушает прямо сейчас.
- Разработка и отладка.
SUBSCRIBEвvalkey-cli, чтобы смотреть, как идут события.
import { createClient } from "redis";const pub = createClient({ url: process.env.VALKEY_URL });const sub = pub.duplicate();await Promise.all([pub.connect(), sub.connect()]);await sub.subscribe("orders:updates", (message) => { broadcastToBrowsers(JSON.parse(message));});await pub.publish("orders:updates", JSON.stringify({ id: 512, status: "shipped" }));Обратите внимание на duplicate(): второе подключение для подписчика по причине, описанной выше.
Уведомления о пространстве ключей
Valkey умеет публиковать и события о собственных ключах: ключ истёк, удалён, записан. Это уведомления о пространстве ключей (keyspace notifications), и они доставляются через обычные каналы pub/sub, например __keyevent@0__:expired. По умолчанию они выключены, потому что их генерация стоит CPU при каждой записи, а включаются настройкой notify-keyspace-events - например, Ex для событий истечения.
Ими соблазнительно воспользоваться для «сделать что-то, когда эта сессия истечёт» или «запустить задачу, когда истечёт ключ-таймер», но они наследуют все слабости pub/sub. Слушатель, отключённый в момент события, никогда о нём не узнает. К тому же события истечения срабатывают, когда Valkey действительно удаляет ключ, а это происходит, когда к ключу обращаются или его находит фоновый цикл истечения, - не обязательно ровно в ту секунду, когда закончился TTL. Для всего, что обязательно должно произойти, надёжная версия той же идеи - sorted set с оценкой по времени выполнения, который опрашивает исполнитель. На размещённом у провайдера экземпляре изменение notify-keyspace-events требует права на CONFIG SET, которого у вашего пользователя может не быть.
Streams: журнал только для добавления#
Stream - это ключ, содержащий упорядоченную последовательность записей. Каждая запись - небольшой набор пар «поле - значение» с id вида <milliseconds>-<sequence>, который назначает сервер.
XADD orders * type created id 512 total 49.90"1791468000123-0"XLEN orders(integer) 1XRANGE orders - +1) 1) "1791468000123-0" 2) 1) "type" 2) "created" 3) "id" 4) "512" 5) "total" 6) "49.90"* просит Valkey сгенерировать id. XRANGE с - и + читает всё от самой старой записи до самой новой; XREVRANGE читает в обратную сторону. Поскольку id - это временные метки, можно прочитать «всё начиная с 09:00», указав в качестве начала значение в миллисекундах.
Без групп потребителей любое число читателей может независимо следить за stream через XREAD, каждый запоминая последний увиденный id:
XREAD COUNT 100 BLOCK 5000 STREAMS orders 1791468000123-0BLOCK 5000 ждёт новых записей до пяти секунд; $ вместо id означает «только записи, добавленные с этого момента». Это уже большой шаг вперёд по сравнению с pub/sub: отключившийся читатель может продолжить со своего последнего id и ничего не пропустить, если записи ещё не обрезаны. Каждый читатель видит каждую запись, так что это рассылка с возможностью повтора.
На этом этапе стоит принять несколько проектных решений. Делайте записи маленькими и самоописательными: тип события, id изменившейся сущности и, возможно, версия, а не вся запись целиком - потребители могут загрузить текущее состояние из базы данных, а маленькие записи делают хранение stream дешёвым. Используйте отдельный stream на каждый вид событий (orders, payments), а не один на всё, чтобы потребителям, которым важны заказы, не приходилось читать и пропускать каждый платёж. И пусть id назначает сервер; задавать свои можно, но они должны всегда возрастать, а при нескольких производителях в этом легко ошибиться.
Группы потребителей: разделение работы с подтверждениями#
Чтобы распределить записи между несколькими исполнителями так, чтобы каждая обрабатывалась один раз, создайте группу потребителей. Группа хранит последний доставленный id для всей группы, а также список ожидающих записей (PEL) - сообщений, доставленных потребителю, но ещё не подтверждённых.
XGROUP CREATE orders billing $ MKSTREAMXREADGROUP GROUP billing worker-1 COUNT 10 BLOCK 5000 STREAMS orders >XACK orders billing 1791468000123-0Что делает каждая из этих трёх строк:
XGROUP CREATE ... $начинает группу с конца stream; используйте0, чтобы обработать всё, что в нём уже есть.MKSTREAMсоздаёт stream, если его ещё нет.XREADGROUPсо специальным id>запрашивает записи, ещё никому в группе не доставленные. Каждая запись уходит ровно одному потребителю. Имя потребителя (worker-1) произвольное, но должно оставаться одинаковым между перезапусками одного и того же исполнителя.XACKудаляет запись из списка ожидающих. До этого группа помнит, что она уworker-1.
Несколько групп могут читать один stream независимо: и billing, и email получают каждый заказ, а внутри каждой группы работа делится между потребителями. Это шаблон для журнала событий, на который реагируют несколько подсистем.
import os, redisr = redis.Redis.from_url(os.environ["VALKEY_URL"], decode_responses=True)GROUP, NAME = "billing", os.environ.get("WORKER_NAME", "worker-1")try: r.xgroup_create("orders", GROUP, id="0", mkstream=True)except redis.ResponseError as e: if "BUSYGROUP" not in str(e): raisewhile True: # first, anything delivered to me before a crash and never acknowledged for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": "0"}, count=50) or []: for entry_id, fields in entries: handle(fields); r.xack("orders", GROUP, entry_id) for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": ">"}, count=50, block=5000) or []: for entry_id, fields in entries: handle(fields); r.xack("orders", GROUP, entry_id)Чтение с id 0 вместо > возвращает собственные ожидающие записи этого потребителя - то, что он получил до падения и так и не подтвердил. Обработка их первым делом при запуске позволяет перезапущенному исполнителю доделать начатое.
Гарантии доставки и список ожидающих записей#
Streams с группами потребителей дают доставку хотя бы один раз. Потребитель, который обработал запись и упал до XACK, увидит её снова; увидит и другой потребитель, который её заберёт. Поэтому обработчики должны быть идемпотентными - безопасными при повторном запуске, - точно как с любой очередью.
Записи, чей потребитель умер навсегда, должен забрать кто-то другой:
XPENDING orders billingXPENDING orders billing - + 10XAUTOCLAIM orders billing worker-2 60000 0-0 COUNT 10XPENDING показывает сводку по незавершённому, а с диапазоном перечисляет записи с владельцем, временем простоя и числом доставок. XAUTOCLAIM передаёт записи, простаивающие дольше 60 000 мс, потребителю worker-2 и возвращает их. Запускайте его периодически в каждом исполнителе, и у вас будет автоматическое восстановление после мёртвых потребителей.
Число доставок - ваш детектор «отравленных» сообщений. Запись, доставленная десять раз, не преуспеет на одиннадцатый; скопируйте её в stream недоставленных через XADD orders:dead ..., подтвердите через XACK в исходной группе и поднимите оповещение. Без этого шага одно некорректное сообщение будет бесконечно ходить по вашим исполнителям.
Удаляйте несуществующих больше потребителей через XGROUP DELCONSUMER, предварительно забрав их ожидающие записи - удаление потребителя отбрасывает его список ожидающих.
Обрезка, память и сохранение данных#
Stream хранит каждую запись, пока вы её не удалите. Подтверждение записи её не удаляет - XACK только обновляет список ожидающих в группе. Без обрезки загруженный stream растёт, пока не закончится память.
XADD orders MAXLEN ~ 100000 * type created id 513XTRIM orders MAXLEN ~ 100000XTRIM orders MINID ~ 1790863200000MAXLEN ограничивает число записей; MINID удаляет всё старше указанного id, то есть задаёт хранение по времени («хранить семь дней»). ~ делает обрезку приблизительной, позволяя Valkey удалять целые внутренние узлы за раз; записей может остаться чуть больше предела, зато это гораздо дешевле точной обрезки. Используйте её, если вам не нужно точное число.
Обрезка не проверяет, обработала ли запись каждая группа. Если группа потребителей отстанет достаточно сильно, обрезка удалит записи, которых она так и не увидела. Задавайте срок хранения с большим запасом относительно того, сколько потребитель может лежать, и следите за отставанием групп - XINFO GROUPS orders показывает для каждой группы последний доставленный id, число ожидающих и, в актуальных версиях, величину отставания.
Память на запись зависит от числа полей и размера значений; маленькие записи стоят по несколько десятков байт, потому что streams хранят их компактно. Сто тысяч небольших событий обычно помещаются в несколько десятков мегабайт - измеряйте через MEMORY USAGE orders, а не доверяйте оценкам. В отличие от pub/sub, streams - это ключи, поэтому AOF и снимки их сохраняют, и политика вытеснения к ним применяется. Stream, который не должен терять записи, требует noeviction, точно как очередь; окно потерь при жёсткой остановке описано в статье сохранение данных в Valkey: AOF и RDB.
В RE:NODE тариф Valkey сохраняет данные на диск через AOF и снимки, так что записи stream и состояние групп потребителей переживают перезапуск. К нему подключаются напрямую по хосту и порту с паролем, сгенерированным для этого сервера; слота прокси перед ним нет, а именно это нужно долгим блокирующим чтениям.
Устранение неполадок#
Подписчики пропускают сообщения. Для pub/sub это ожидаемо, когда они отключены, медленные или ещё не подписались. Если это неприемлемо, архитектуре нужен stream.
`PUBLISH` возвращает 0. В этот момент у канала нет подписчиков. Проверьте имя канала - каналы pub/sub чувствительны к регистру и никак не связаны с номерами логических баз.
Отставание группы потребителей постоянно растёт. Потребители медленнее производителей. Добавьте потребителей в группу, увеличьте COUNT или найдите медленный обработчик.
Список ожидающих растёт без предела. Обработчики не вызывают XACK или падают до него. Посмотрите в XPENDING владельцев и число доставок.
`BUSYGROUP Consumer Group name already exists`. При запуске это безвредно; перехватывайте ошибку, как в примере.
Память равномерно растёт вместе со stream. Нет обрезки. Добавьте MAXLEN ~ к XADD или запускайте XTRIM по расписанию.
FAQ#
Можно ли получить сообщения pub/sub, отправленные, пока приложение лежало?
Нет. Pub/sub ничего не хранит. Если это нужно, публикуйте в stream, и пусть каждый читатель отслеживает свой последний id; читатели, которые лежали, продолжат с того места, где остановились.
Заменяет ли stream в Valkey Kafka?
Для одного приложения или нескольких служб он покрывает тот же шаблон - упорядоченный журнал с независимыми группами потребителей - при гораздо меньших затратах на эксплуатацию. Kafka создана для хранения, измеряемого терабайтами, разбиения на партиции по многим машинам и повторного проигрывания за недели. Если ваш stream спокойно помещается в память, Valkey - соразмерный инструмент.
Что использовать: streams или библиотеку очередей задач?
Библиотеку очередей задач - когда сообщения являются задачами с повторами, задержками и расписанием: BullMQ, Celery и RQ всё это уже дают. Streams - когда сообщения являются событиями, на которые реагируют несколько потребителей, или когда нужна упорядоченная история, которую можно проиграть заново.
Как сделать обработку ровно один раз?
Одним только транспортом - никак; streams доставляют хотя бы один раз. Сделайте обработчик идемпотентным: записывайте id записи или бизнес-ключ в базу данных в той же транзакции, что и саму работу, и пропускайте всё уже записанное.
Учитываются ли сообщения pub/sub в лимитах памяти?
Только пока они буферизуются для медленных подписчиков, а эти буферы ограничены для каждого клиента настройкой client-output-buffer-limit. Записи streams учитываются полностью, потому что это хранимые данные.




Комментарии
Полностью анонимно: без аккаунта, без почты, без cookie. Мы храним имя, которое вы ввели, текст и время - больше ничего. Количество ссылок ограничено, разметка не отображается.