Redis · Стримы

Стримы

Стрим — журнал с добавлением в конец: у каждой записи есть идентификатор, присвоенный сервером, и набор полей. От списка он отличается тем, что чтение не забирает запись, а группы потребителей заставляют сам Redis вести учёт: кто что взял и что ещё не подтверждено.

Что такое стрим

Запись выглядит так: идентификатор 1755600000123-0 (миллисекунды серверных часов и порядковый номер внутри этой миллисекунды) и поля {"type": "signup", "user": "42"}. Идентификаторы монотонно растут, поэтому «прочитать всё после вот этого» — обычная операция, а не поиск.

Свойства, из которых следует всё остальное:

  • Чтение ничего не изымает. Запись живёт, пока её не удалят или не подрежут журнал. Одну и ту же запись могут прочитать сколько угодно независимых читателей.
  • Порядок гарантирован и совпадает с порядком добавления.
  • Журнал растёт вечно, если его не ограничить. Это главная плата за историю.
  • Группы ведёт сервер. Он помнит, какая запись какому потребителю ушла и была ли подтверждена. Список этого не умеет — там учёт приходится строить руками.

Чем это отличается от списка

Список Стрим
После чтения элемент исчез запись на месте
Читателей на сообщение один сколько угодно
История нет есть, пока не подрезали
Учёт «взято, но не завершено» вручную (moveTo + remove) встроенный, в группе
Рост ограничен потреблением вечный, пока не ограничить
Идентификаторы нет есть, по ним читают и подтверждают

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

Где принято применять

Шина событий. Одно событие — несколько независимых обработчиков: письмо, метрика, вебхук. Каждый читает своей группой и не мешает остальным.

Очередь с гарантией. Воркер взял запись, упал — она осталась в списке незавершённых и достанется другому. Ровно то, что в списках приходится собирать из moveTo() и remove().

Журнал изменений. Аудит, история заказа, лента активности: можно перечитать, отмотать назад, разобрать инцидент по следам.

Буфер с догоняющим чтением. Обработчик пишет, потребитель читает в своём темпе и после перезапуска продолжает с того места, где остановился.

Где стрим — неправильный выбор:

Задача Почему не стрим Что вместо
Простая очередь без истории лишний рост и лишние понятия список
Уведомить всех подключённых сейчас стрим хранит, а не рассылает pub/sub
Найти запись по содержимому поиска нет, только по идентификатору и диапазону индекс на множествах
Хранить состояние объекта стрим про события, а не про текущее значение хеш

Ручка

php
$events = $store->stream('events');    // сервер увидит 'session:events'

stream() ничего не выполняет — это представление, а не запрос. Префикс применяется здесь один раз. Соединения ручка не держит, кроме одного случая: follow() открывает собственное, потому что блокирующее чтение занимает соединение на всё время ожидания.

Записи

Методы чтения возвращают StreamEntry:

php
foreach ($events->range(count: 100) as $entry) {
  $entry->id;                  // '1755600000123-0'
  $entry->fields;              // ['type' => 'signup', 'user' => '42']
  $entry->get('type');         // 'signup'
  $entry->get('absent', '');  // значение по умолчанию
  $entry->has('user');         // bool
  $entry->timestamp();         // 1755600000123 — миллисекунды серверных часов
}

Драйвер отдаёт записи как ['1755600000123-0' => ['type' => 'signup']] — массив, у которого единственный ключ и есть идентификатор. Обходить такое приходится через key()/current(), а идентификатор нужен постоянно: им подтверждают обработку и запоминают позицию.


Справочник RedisStream

Запись

add()

php
public function add(array $fields, ?int $cap = null, bool $exact = false, string $id = '*'): string

Добавляет запись в конец журнала.

Аргумент Тип По умолчанию Что делает
$fields array поля записи. Хотя бы одно: пустых записей в Redis нет
$cap ?int null ограничить журнал примерно этим числом записей
$exact bool false подрезать ровно до $cap, а не приблизительно
$id string '*' идентификатор; * — сервер назначит сам

Возвращает идентификатор, присвоенный записи. Бросает RedisCommandException, если ключ занят структурой другого типа.

php
$events->add(['type' => 'signup', 'user' => '42']);
// '1755600000123-0'

$events->add(['level' => 'warn', 'msg' => $text], cap: 10_000);

Свои идентификаторы почти никогда не нужны: серверные монотонны и включают время. Ручное назначение имеет смысл только при переносе данных из другого журнала.

Значения полей — строки

Как и везде в Redis. Массив в поле превратится в "Array", если у конфигурации не задан сериализатор.

trim()

php
public function trim(int $cap, bool $exact = false): int

Оставляет примерно $cap новейших записей.

Аргумент Тип По умолчанию Что делает
$cap int сколько записей оставить
$exact bool false подрезать ровно, а не приблизительно

Возвращает число удалённых записей.

Приблизительная подрезка может не удалить ничего

Она останавливается на границе внутреннего блока стрима — так дешевле, и Redis рекомендует именно её. Следствие неочевидное: на коротком журнале trim(10) вполне может вернуть 0 и оставить всё как было. Это не сбой; если нужна предсказуемая длина, просите exact: true и платите за обход.

php
$events->trim(10_000);                 // дёшево, длина «примерно такая»
$events->trim(10_000, exact: true);    // ровно столько, дороже

trimBefore()

php
public function trimBefore(string $id, bool $exact = false): int

Удаляет всё, что старше указанного идентификатора.

Аргумент Тип По умолчанию Что делает
$id string граница; записи с меньшим идентификатором удаляются
$exact bool false точная граница вместо приблизительной

Возвращает число удалённых записей.

Так журнал ограничивают по времени, а не по количеству: идентификатор начинается с миллисекунд серверных часов, поэтому «старше суток» — это идентификатор:

php
$events->trimBefore((string) ((time() - 86400) * 1000));

delete()

php
public function delete(string ...$ids): int

Удаляет записи по идентификаторам. Возвращает число существовавших. Без аргументов — 0.

Удаление из середины журнала допустимо, но это исключение: обычный способ ограничить стрим — подрезка. delete() нужен, когда конкретную запись нельзя больше показывать.

Чтение

count()

php
public function count(): int

Число записей; для несуществующего ключа — 0. Аргументов нет.

range()

php
public function range(?string $from = null, ?string $to = null, ?int $count = null): array

Записи в порядке добавления, ничего не изымая.

Аргумент Тип По умолчанию Что делает
$from ?string null с какого идентификатора, включительно; null — с начала
$to ?string null по какой, включительно; null — до конца
$count ?int null не больше стольких записей

Возвращает list<StreamEntry>; для пустого журнала — пустой массив.

php
$events->range();                    // всё
$events->range(count: 100);          // первая сотня
$events->range($since, count: 50);   // с известной позиции

reverse()

php
public function reverse(?int $count = null): array

То же самое, но новейшие первыми — дешёвый способ посмотреть, что только что произошло.

php
$events->reverse(count: 10);   // последние десять, свежие сверху

after()

php
public function after(string $id, int $count = 10): array

Записи, добавленные после указанного идентификатора, без ожидания.

Аргумент Тип По умолчанию Что делает
$id string позиция, после которой читать
$count int 10 не больше стольких записей

Возвращает list<StreamEntry>; пустой массив, если ничего нового нет.

Это основа читателя, который сам хранит позицию: запомнили идентификатор последней обработанной записи, передали его в следующий раз. Два таких читателя не мешают друг другу — стрим чтением не расходуется.

php
$position = $this->positions->get('events') ?? '0';

foreach ($events->after($position, count: 100) as $entry) {
  $this->handle($entry);
  $position = $entry->id;
}

$this->positions->set('events', $position);

follow()

php
public function follow(float $timeout = 0.0, int $count = 10, ?string $from = null): array

Ждёт записи новее тех, что эта ручка уже видела.

Аргумент Тип По умолчанию Что делает
$timeout float 0.0 сколько секунд ждать; 0 — неограниченно
$count int 10 не больше стольких записей за раз
$from ?string null позиция для первого вызова; '0' — начать со всей истории

Возвращает list<StreamEntry>; пустой массив, если за отведённое время ничего не появилось.

main/Daemons/EventTail.php
$events = $store->stream('events');

while (!$this->stopping) {
  foreach ($events->follow(timeout: 5) as $entry) {
      $this->handle($entry);
  }
}

$events->close();

Позиция запоминается — и почему это не «$»

Первый вызов начинает с текущего конца журнала и тратит на это один лишний запрос: спрашивает идентификатор последней записи. Можно было бы обойтись редисовым $ («записи, добавленные после начала этого вызова»), но между двумя вызовами цикла есть промежуток, и всё записанное в нём не попало бы ни в один — работа терялась бы без следа. Разрешив позицию один раз, ручка превращает хвост в непрерывную цепочку идентификаторов.

История через follow() не читается намеренно: для неё есть range() и after().

follow() не берёт соединение из пула

Блокирующее чтение занимает соединение на всё время ожидания, поэтому ручка открывает своё. close() его возвращает. По этой же причине ручку стоит держать на весь цикл, а не брать заново на каждой итерации.

Наблюдение

info()

php
public function info(): array

Сводка сервера о журнале: length, first-entry, last-entry, last-generated-id, число групп. Аргументов нет.

php
$events->info()['length'];              // 1204
$events->info()['last-generated-id'];   // '1755600000123-0'

groups()

php
public function groups(): array

Все группы потребителей на этом стриме, как их видит сервер: name, consumers, pending, last-delivered-id, lag.

php
foreach ($events->groups() as $group) {
  if ($group['pending'] > 1000) {
      $this->alert("группа {$group['name']} не успевает");
  }
}

lag — сколько записей группа ещё не видела; вместе с pending это два разных отставания: первое про «не прочитано», второе про «прочитано, но не подтверждено».

Ключ целиком

php
$events->expireKey(86400);   // срок жизни всему журналу
$events->keyTtl();           // 86400 | null
$events->deleteKey();        // удалить журнал вместе с группами
$events->name();             // 'session:events'
$events->close();            // закрыть соединение, открытое follow()

deleteKey() уносит и записи, и группы с их учётом. Отдельных сроков жизни у записей нет — ограничивают подрезкой.


Группы потребителей

Зачем они

Без группы каждый читатель видит все записи и сам следит за своей позицией. Группа меняет обе части: запись уходит одному потребителю группы, и сервер держит её в списке незавершённых, пока её не подтвердят.

Отсюда гарантия, ради которой это и берут: воркер взял запись и умер — она не потеряна. Она числится за ним, её видно в pending(), и другой воркер может её забрать.

php
$events->ensureGroup('mailers');                          // один раз, при старте
$group = $events->group('mailers', consumer: 'worker-1');

ensureGroup()

php
public function ensureGroup(string $name, string $from = '0'): bool
Аргумент Тип По умолчанию Что делает
$name string имя группы
$from string '0' откуда группа начинает читать: '0' — с начала журнала, '$' — только новое

Возвращает true, если группа создана, и false, если она уже была — второе не ошибка. Стрим при необходимости создаётся тоже, поэтому группу можно объявить на старте, до первой записи.

Почему группа не создаётся сама при consume()

Тогда опечатка в имени превращалась бы в новую пустую группу, которая молча ничего не получает. Явное создание оставляет опечатке единственный исход — ошибку NOGROUP, называющую имя.

group()

php
public function group(string $name, string $consumer = 'default'): RedisStreamGroup

Ручка на группу от лица конкретного потребителя. Ничего не создаёт и не отправляет.


Справочник RedisStreamGroup

Потребление

consume()

php
public function consume(int $count = 10, float $timeout = 0.0): array

Забирает записи, которые ещё никому в группе не выдавались.

Аргумент Тип По умолчанию Что делает
$count int 10 не больше стольких записей
$timeout float 0.0 секунд ожидания; 0 — неограниченно, отрицательное — не ждать вовсе

Возвращает list<StreamEntry>. Бросает RedisCommandException с NOGROUP, если группы не существует.

Каждая запись достаётся ровно одному потребителю группы и попадает в его список незавершённых, пока не будет подтверждена.

php
$group->consume(count: 10, timeout: 5);   // ждёт, своё соединение
$group->consume(timeout: -1);             // заглянуть и вернуться, соединение из пула

Отрицательный таймаут — единственный режим, который не берёт отдельное соединение; он для проверок и разовых заглядываний.

backlog()

php
public function backlog(int $count = 10): array

Возвращает то, что этот потребитель уже взял, но не подтвердил.

Возвращает list<StreamEntry>, не ждёт.

Это первое, что должен сделать перезапустившийся воркер: записи, взятые до перезапуска, числятся за ним и через consume() больше не придут. Прочитать их обратно — способ продолжить работу, не дожидаясь, пока их кто-то перехватит по простою.

main/Daemons/MailWorker.php
$group = $events->group('mailers', consumer: $this->workerName());

// 1. доделать своё
foreach ($group->backlog() as $entry) {
  $this->handle($entry);
  $group->ack($entry);
}

// 2. и только потом брать новое
while (!$this->stopping) {
  foreach ($group->consume(count: 10, timeout: 5) as $entry) {
      $this->handle($entry);
      $group->ack($entry);
  }
}

$group->close();

ack()

php
public function ack(StreamEntry|string ...$entries): int

Отмечает записи обработанными, убирая их из списка незавершённых.

Аргумент Тип По умолчанию Что делает
...$entries StreamEntry|string записи или их идентификаторы

Возвращает число записей, которые числились незавершёнными и теперь подтверждены. Без аргументов — 0.

php
$group->ack($entry);              // объект
$group->ack($entry->id);          // или идентификатор
$group->ack(...$entries);         // или пачкой

Подтверждение не удаляет запись из журнала

ack() закрывает учёт в группе: запись перестаёт числиться незавершённой. Сама запись остаётся в стриме и по-прежнему видна в range() — на то он и журнал. Ограничивать рост нужно отдельно, подрезкой.

Подтверждать автоматически, при выдаче, пакет не станет: это уничтожило бы единственную гарантию группы — что запись переживёт воркера, который её взял и не закончил.

Восстановление

pending()

php
public function pending(int $count = 100, ?string $consumer = null): array

Записи, выданные группой и никем не подтверждённые.

Аргумент Тип По умолчанию Что делает
$count int 100 не больше стольких
$consumer ?string null ограничить одним потребителем; null — по всей группе

Возвращает list<PendingEntry>:

php
foreach ($group->pending() as $stuck) {
  $stuck->id;           // '1755600000123-0'
  $stuck->consumer;     // 'worker-1'
  $stuck->idleMs;       // 61240 — сколько миллисекунд назад её выдали
  $stuck->deliveries;   // 3 — сколько раз выдавали
}

Две цифры отвечают на разные вопросы. Большой idleMs при одной доставке — потребитель умер, держа запись. Растущий deliveries — запись роняет всех, кто за неё берётся.

Что делать со второй — перезапустить, разбудить алерт, отложить в «мёртвые» — решает приложение: цену повтора знает только автор задачи.

pendingCount()

php
public function pendingCount(): int

Сколько всего записей группа выдала и не получила подтверждения. Аргументов нет.

claimStale()

php
public function claimStale(int $idle, int $count = 10, string $from = '0-0'): array

Забирает себе записи, слишком долго висящие у другого потребителя.

Аргумент Тип По умолчанию Что делает
$idle int сколько миллисекунд запись должна простаивать, чтобы её можно было забрать
$count int 10 не больше стольких
$from string '0-0' с какого места просматривать список незавершённых

Возвращает list<StreamEntry> — уже забранные записи, готовые к обработке.

php
foreach ($group->claimStale(idle: 60_000) as $entry) {
  $this->handle($entry);
  $group->ack($entry);
}

Повторный вызов не вернёт те же записи

Захват сбрасывает счётчик простоя, поэтому следующий вызов с тем же порогом их уже не увидит — цикл сходится, а не крутится на одном месте. Счётчик доставок при этом растёт, и по нему видно, что запись пошла по второму кругу.

Порог подбирают по самой долгой честной обработке: если письмо отправляется до тридцати секунд, idle: 60_000 не отберёт работу у живого воркера.

consumers()

php
public function consumers(): array

Все потребители группы, как их видит сервер: name, pending, idle. Аргументов нет.

php
foreach ($group->consumers() as $consumer) {
  if ($consumer['idle'] > 300_000 && $consumer['pending'] > 0) {
      $this->alert("{$consumer['name']} молчит пять минут и держит работу");
  }
}

destroy()

php
public function destroy(): bool

Удаляет группу вместе с её учётом и потребителями. Записи журнала остаются.

Служебное

php
$group->as('worker-2');    // та же группа от лица другого потребителя
$group->name();            // 'mailers'
$group->consumer();        // 'worker-1'
$group->close();           // закрыть соединение, открытое блокирующим consume()

Имена потребителей должны быть уникальны

Сервер ведёт список незавершённых на потребителя. Два воркера с одним именем делят один список и будут разбирать работу друг друга — включая ту, что первый прямо сейчас обрабатывает. Называйте по чему-то устойчивому и различимому: слот воркера, хост и pid.

Соответствие командам Redis

Метод Команда
add() XADD, с capXADD ... MAXLEN
trim() / trimBefore() XTRIM MAXLEN / XTRIM MINID
delete() / count() XDEL / XLEN
range() / reverse() XRANGE / XREVRANGE
after() / follow() XREAD / XREAD BLOCK
info() / groups() / consumers() XINFO STREAM / GROUPS / CONSUMERS
ensureGroup() / destroy() XGROUP CREATE / XGROUP DESTROY
consume() / backlog() XREADGROUP с > / с 0
ack() XACK
pending() / pendingCount() XPENDING
claimStale() XAUTOCLAIM

Чего здесь нет — XCLAIM поштучно, XSETID, XGROUP CREATECONSUMER — доступно через raw() вместе с name().

Дальше

  • Списки — очередь, когда история не нужна
  • Хеши — состояние объекта рядом с журналом событий
  • Пул соединений — почему блокирующие чтения держатся в стороне от пула