Redis · Стримы

Стримы

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

Пакет flytachi/winter-redisКлассы RedisStream · RedisStreamGroupЗаписи StreamEntry · PendingEntry

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

Проблема. Очередь на списке отвечает на один вопрос — «что делать дальше» — и теряет всё остальное. Прочитанное исчезает: нельзя перечитать событие, нельзя дать одно событие двум разным обработчикам, нельзя после инцидента посмотреть, что вообще происходило. А если воркер взял задачу и умер, задача исчезает вместе с ним, и узнать об этом неоткуда: учёт «взято, но не завершено» приходится строить вручную.

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

php
$events = $store->stream('events');

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

$events->ensureGroup('mailers');
$group = $events->group('mailers', consumer: 'worker-1');

foreach ($group->consume(count: 10, timeout: 5) as $entry) {
  $this->handle($entry);
  $group->ack($entry);        // пока не подтвердили — запись числится за воркером
}

Свойства структуры

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

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

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

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

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

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

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

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

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

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

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

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

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

Записи: StreamEntry

Методы чтения возвращают не сырые массивы, а объекты StreamEntry.

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

Справочник RedisStream

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

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

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

Запись

add()

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

Синтаксис

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

Параметры

$fields — поля записи. Хотя бы одно: пустых записей в Redis не бывает. Значения — строки, как и везде.

$cap — ограничить журнал примерно этим числом записей прямо при добавлении. По умолчанию null — не ограничивать.

$exact — подрезать ровно до $cap, а не приблизительно. По умолчанию false; разница разобрана в trim().

$id — идентификатор записи. По умолчанию '*' — сервер назначит сам. Свой идентификатор обязан быть больше последнего в журнале.

Возвращает

Идентификатор, присвоенный записи.

Ошибки

RedisCommandException — если ключ занят структурой другого типа, либо если переданный $id не больше последнего:

text
RedisCommandException: ERR The ID specified in XADD is equal or smaller
than the target stream top item

Пример

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

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

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

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

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

trim()

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

Синтаксис

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

Параметры

$cap — сколько записей оставить.

$exact — подрезать ровно, а не приблизительно. По умолчанию false.

Возвращает

Число удалённых записей.

Пример

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

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

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

Это не сбой. Если нужна предсказуемая длина, просите exact: true и платите за обход.

trimBefore()

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

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

Синтаксис

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

Параметры

$id — граница; записи с меньшим идентификатором удаляются.

$exact — точная граница вместо приблизительной. По умолчанию false, и оговорка из trim() действует здесь точно так же.

Возвращает

Число удалённых записей.

Пример

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

delete()

Удаляет записи по идентификаторам.

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

Синтаксис

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

Параметры

...$ids — идентификаторы записей, сколько угодно. Можно не передать ни одного.

Возвращает

Число записей, которые существовали и были удалены. Без аргументов — 0.

Чтение

count()

Сообщает число записей в журнале.

Синтаксис

php
public function count(): int

Параметры

Нет.

Возвращает

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

range()

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

Синтаксис

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

Параметры

$from — с какого идентификатора читать, включительно. По умолчанию null — с начала журнала.

$to — по какой идентификатор, включительно. По умолчанию null — до конца.

$count — не больше стольких записей. По умолчанию null — без ограничения.

Возвращает

Массив StreamEntry; для пустого журнала — пустой массив.

Пример

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

reverse()

Возвращает записи новейшими вперёд.

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

Синтаксис

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

Параметры

$count — не больше стольких записей. По умолчанию null — весь журнал.

Возвращает

Массив StreamEntry, свежие первыми.

Пример

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

after()

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

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

Синтаксис

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

Параметры

$id — позиция, после которой читать. Сама запись с этим идентификатором в результат не попадёт.

$count — не больше стольких записей. По умолчанию 10.

Возвращает

Массив 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 — сколько секунд ждать. По умолчанию 0.0 — неограниченно.

$count — не больше стольких записей за один вызов. По умолчанию 10.

$from — позиция для первого вызова. По умолчанию null — начать с текущего конца журнала; '0' — прочитать сначала всю историю.

Возвращает

Массив 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 это два разных отставания: первое про «не прочитано», второе про «прочитано, но не подтверждено».

Ключ целиком

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

expireKey()

Задаёт срок жизни всему журналу.

Синтаксис

php
public function expireKey(int $seconds): bool

Параметры

$seconds — через сколько секунд удалить ключ целиком, считая от текущего момента.

Возвращает

true, если срок установлен; false, если ключа нет.

keyTtl()

Сообщает, сколько журналу осталось жить.

Синтаксис

php
public function keyTtl(): ?int

Параметры

Нет.

Возвращает

Оставшиеся секунды, либо null — и когда срока нет, и когда нет ключа.

deleteKey()

Удаляет журнал вместе с группами.

Синтаксис

php
public function deleteKey(): bool

Параметры

Нет.

Возвращает

true, если ключ существовал.

Уносит и записи, и группы с их учётом незавершённого.

name()

Возвращает имя ключа с префиксом — то, что видит сервер.

Синтаксис

php
public function name(): string

Параметры

Нет.

Возвращает

Полное имя ключа. Нужно для команд, которых нет в ручке: их выполняют через raw().

close()

Закрывает соединение, которое открыл follow().

Вызывать безопасно всегда: если блокирующего чтения не было, метод ничего не делает.

Синтаксис

php
public function close(): void

Параметры

Нет.

Возвращает

Ничего.


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

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

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

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

Создание

ensureGroup()

Создаёт группу потребителей, если её ещё нет.

Стрим при необходимости создаётся тоже, поэтому группу можно объявить на старте приложения, до первой записи.

Синтаксис

php
public function ensureGroup(string $name, string $from = '0'): bool

Параметры

$name — имя группы.

$from — откуда группа начинает читать. По умолчанию '0' — с начала журнала, то есть группа получит всю накопленную историю. '$' — только записи, добавленные после создания группы.

Возвращает

true, если группа создана; false, если она уже была — второе не ошибка, и повторный вызов при каждом старте безопасен.

Пример

php
$events->ensureGroup('mailers');          // с начала журнала
$events->ensureGroup('metrics', '$');     // только новое

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

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

group()

Возвращает ручку на группу от лица конкретного потребителя.

Ничего не создаёт и на сервер не ходит — как и остальные ручки, это представление.

Синтаксис

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

Параметры

$name — имя группы. Она должна существовать: создаёт её ensureGroup().

$consumer — имя потребителя, от лица которого будут идти команды. По умолчанию 'default'.

Возвращает

RedisStreamGroup.

Пример

php
$group = $events->group('mailers', consumer: 'worker-1');

Справочник RedisStreamGroup

RedisStreamGroup — ручка на пару «группа + потребитель». Имя потребителя в ней не случайно: сервер ведёт список незавершённых записей на потребителя, а не на группу, поэтому все методы восстановления работают от чьего-то лица.

Потребление

consume()

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

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

Синтаксис

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

Параметры

$count — не больше стольких записей за вызов. По умолчанию 10.

$timeout — сколько секунд ждать новых записей. По умолчанию 0.0 — ждать неограниченно. Отрицательное значение означает «не ждать вовсе»: заглянуть и сразу вернуться.

Возвращает

Массив StreamEntry; пустой, если новых записей нет.

Ошибки

RedisCommandException с текстом NOGROUP — если группы не существует.

Пример

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

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

backlog()

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

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

Синтаксис

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

Параметры

$count — не больше стольких записей. По умолчанию 10.

Возвращает

Массив StreamEntry. Не ждёт: если незавершённого нет, сразу пустой массив.

Пример

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 — записи либо их идентификаторы, вперемешку и сколько угодно. Можно не передать ни одной.

Возвращает

Число записей, которые числились незавершёнными и теперь подтверждены. Без аргументов — 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 — не больше стольких записей. По умолчанию 100.

$consumer — ограничить одним потребителем. По умолчанию null — по всей группе.

Возвращает

Массив 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 — сколько миллисекунд запись должна простаивать, чтобы её можно было забрать. Порог подбирают по самой долгой честной обработке: если письмо отправляется до тридцати секунд, idle: 60_000 не отберёт работу у живого воркера.

$count — не больше стольких записей. По умолчанию 10.

$from — с какого места просматривать список незавершённых. По умолчанию '0-0' — с начала.

Возвращает

Массив StreamEntry — уже забранные записи, готовые к обработке.

Пример

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

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

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

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

Параметры

Нет.

Возвращает

true, если группа существовала.

Служебное

as()

Возвращает ту же группу от лица другого потребителя.

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

Синтаксис

php
public function as(string $consumer): RedisStreamGroup

Параметры

$consumer — имя потребителя.

Возвращает

Новую RedisStreamGroup с тем же стримом и группой.

Пример

php
$group->as('worker-2')->backlog();    // что не доделал второй воркер

RedisStreamGroup::name()

Возвращает имя группы.

Синтаксис

php
public function name(): string

Возвращает

Имя группы — то, что передали в group().

consumer()

Возвращает имя потребителя, от лица которого работает ручка.

Синтаксис

php
public function consumer(): string

Возвращает

Имя потребителя.

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

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

RedisStreamGroup::close()

Закрывает соединение, открытое блокирующим consume().

Синтаксис

php
public function close(): void

Возвращает

Ничего.

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

Метод Команда
add() XADD, с cap — XADD ... 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().

Дальше

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