Redis Streams давно перестали быть экзотикой и стали нормальным способом передачи событий между сервисами. В PHP есть два популярных подхода для работы с ними: Amp с неблокирующим I/O и Swoole с корутинами. Оба подхода позволяют реализовать устойчивые consumer-группы, ручной ack, автоматическое перенаправление зависших сообщений, backpressure, экспоненциальные ретраи и дед-лейтер.
🛠️ Что строим
Задача — создать шину событий заказов. Продюсер записывает события в
orders:events с помощью команды XADD с триммингом. Несколько воркеров читают из consumer-группы orders:cg с использованием XREADGROUP в блокирующем режиме, подтверждают обработку через XACK, а зависшие записи перенаправляются на активного потребителя через XAUTOCLAIM. Если событие стабильно не обрабатывается, оно отправляется в orders:events:dlq и больше не участвует в основном потоке. Мониторинг задержки группы осуществляется через XINFO GROUPS, а хвосты периодически очищаются.⚙️ Варианты реализации
Amp: Неблокирующее чтение с использованием
XREADGROUP, ограничение параллельной обработки с помощью LocalSemaphore, подтверждение XACK пачками. Для подбора зависших сообщений параллельно запускается цикл XAUTOCLAIM.Swoole: Использование корутин с каналом как семафором, Redis из
ext-phpredis. Параллельная обработка ограничивается размером канала.🔁 Экспоненциальные ретраи
Redis Streams не поддерживают отложенные сообщения, но можно реализовать экспоненциальные ретраи с использованием ZSET. При ошибке вы подтверждаете задачу и кладёте её в
ZSET orders:retry со значением score = now + backoffMs. Отдельная корутина периодически извлекает задачи из ZSET и повторно добавляет их в основной стрим с увеличенным счётчиком попыток.📊 Мониторинг и масштабирование
Для мониторинга используйте команду
XINFO GROUPS, чтобы отслеживать количество записей, которые ещё не доставлены группе. Если лаг стабильно растёт, добавляйте консьюмеров. Если лаг «пилит» около нуля, можно уменьшить число воркеров.🔗 Хабр
Библиотека пхпшника