Стримы
Стрим — журнал с добавлением в конец: у каждой записи есть идентификатор, присвоенный сервером, и набор полей. От списка он отличается тем, что чтение не забирает запись, а группы потребителей заставляют сам Redis вести учёт: кто что взял и что ещё не подтверждено.
Что такое стрим и зачем
Проблема. Очередь на списке отвечает на один вопрос — «что делать дальше» — и теряет всё остальное. Прочитанное исчезает: нельзя перечитать событие, нельзя дать одно событие двум разным обработчикам, нельзя после инцидента посмотреть, что вообще происходило. А если воркер взял задачу и умер, задача исчезает вместе с ним, и узнать об этом неоткуда: учёт «взято, но не завершено» приходится строить вручную.
Решение. Стрим — журнал, из которого чтение ничего не изымает. Запись остаётся на месте, у неё есть идентификатор, и по нему можно вернуться назад или продолжить с того места, где остановились. А группы потребителей перекладывают учёт доставки на сам Redis: он помнит, кому какая запись ушла и подтвердил ли тот её обработку.
$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 кладёт его рядом с полями.
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 — ручка на один стрим: объект, через который выполняются команды над этим
ключом. Получают её у стора:
$events = $store->stream('events'); // сервер увидит 'session:events'stream() ничего не выполняет — это представление, а
не запрос. Префикс применяется здесь один раз. Соединения ручка не держит, кроме одного
случая: follow() открывает собственное, потому что блокирующее чтение
занимает соединение на всё время ожидания.
Запись
add()
Добавляет запись в конец журнала.
Синтаксис
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 не больше последнего:
RedisCommandException: ERR The ID specified in XADD is equal or smaller
than the target stream top itemПример
$events->add(['type' => 'signup', 'user' => '42']);
// '1755600000123-0'
$events->add(['level' => 'warn', 'msg' => $text], cap: 10_000);Свои идентификаторы почти никогда не нужны: серверные монотонны и включают время. Ручное назначение имеет смысл только при переносе данных из другого журнала.
Значения полей — строки
Как и везде в Redis. Массив в поле превратится в "Array", если у конфигурации не задан
сериализатор.
trim()
Оставляет примерно $cap новейших записей.
Синтаксис
public function trim(int $cap, bool $exact = false): intПараметры
$cap — сколько записей оставить.
$exact — подрезать ровно, а не приблизительно. По умолчанию false.
Возвращает
Число удалённых записей.
Пример
$events->trim(10_000); // дёшево, длина «примерно такая»
$events->trim(10_000, exact: true); // ровно столько, дорожеПриблизительная подрезка может не удалить ничего
Она останавливается на границе внутреннего блока стрима — так дешевле, и Redis
рекомендует именно её. Следствие неочевидное и проверено на живом сервере: на журнале из
десяти записей trim(5) возвращает 0 и оставляет все десять, а trim(5, exact: true)
удаляет пять и оставляет пять.
Это не сбой. Если нужна предсказуемая длина, просите exact: true и платите за обход.
trimBefore()
Удаляет всё, что старше указанного идентификатора.
Так журнал ограничивают по времени, а не по количеству: идентификатор начинается с миллисекунд серверных часов, поэтому «старше суток» — это просто идентификатор.
Синтаксис
public function trimBefore(string $id, bool $exact = false): intПараметры
$id — граница; записи с меньшим идентификатором удаляются.
$exact — точная граница вместо приблизительной. По умолчанию false, и оговорка из
trim() действует здесь точно так же.
Возвращает
Число удалённых записей.
Пример
$events->trimBefore((string) ((time() - 86400) * 1000));delete()
Удаляет записи по идентификаторам.
Удаление из середины журнала допустимо, но это исключение: обычный способ ограничить
стрим — подрезка. delete() нужен, когда конкретную запись нельзя больше показывать.
Синтаксис
public function delete(string ...$ids): intПараметры
...$ids — идентификаторы записей, сколько угодно. Можно не передать ни одного.
Возвращает
Число записей, которые существовали и были удалены. Без аргументов — 0.
Чтение
count()
Сообщает число записей в журнале.
Синтаксис
public function count(): intПараметры
Нет.
Возвращает
Количество записей; для несуществующего ключа — 0.
range()
Возвращает записи в порядке добавления, ничего не изымая.
Синтаксис
public function range(?string $from = null, ?string $to = null, ?int $count = null): arrayПараметры
$from — с какого идентификатора читать, включительно. По умолчанию null — с начала
журнала.
$to — по какой идентификатор, включительно. По умолчанию null — до конца.
$count — не больше стольких записей. По умолчанию null — без ограничения.
Возвращает
Массив StreamEntry; для пустого журнала — пустой массив.
Пример
$events->range(); // всё
$events->range(count: 100); // первая сотня
$events->range($since, count: 50); // с известной позицииreverse()
Возвращает записи новейшими вперёд.
Дешёвый способ посмотреть, что только что произошло, не читая журнал с начала.
Синтаксис
public function reverse(?int $count = null): arrayПараметры
$count — не больше стольких записей. По умолчанию null — весь журнал.
Возвращает
Массив StreamEntry, свежие первыми.
Пример
$events->reverse(count: 10); // последние десять, свежие сверхуafter()
Возвращает записи, добавленные после указанного идентификатора, не ожидая новых.
Это основа читателя, который сам хранит позицию: запомнили идентификатор последней обработанной записи, передали его в следующий раз. Два таких читателя не мешают друг другу — стрим чтением не расходуется.
Синтаксис
public function after(string $id, int $count = 10): arrayПараметры
$id — позиция, после которой читать. Сама запись с этим идентификатором в результат не
попадёт.
$count — не больше стольких записей. По умолчанию 10.
Возвращает
Массив StreamEntry; пустой, если ничего нового нет.
Пример
$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()
Ждёт записи новее тех, что эта ручка уже видела.
Синтаксис
public function follow(float $timeout = 0.0, int $count = 10, ?string $from = null): arrayПараметры
$timeout — сколько секунд ждать. По умолчанию 0.0 — неограниченно.
$count — не больше стольких записей за один вызов. По умолчанию 10.
$from — позиция для первого вызова. По умолчанию null — начать с текущего конца
журнала; '0' — прочитать сначала всю историю.
Возвращает
Массив StreamEntry; пустой, если за отведённое время ничего не
появилось.
Пример
$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()
Возвращает сводку сервера о журнале.
Синтаксис
public function info(): arrayПараметры
Нет.
Возвращает
Массив с полями, которые отдаёт сервер: length, first-entry, last-entry,
last-generated-id, число групп и прочее. Для несуществующего ключа — пустой массив.
Пример
$events->info()['length']; // 1204
$events->info()['last-generated-id']; // '1755600000123-0'groups()
Возвращает все группы потребителей на этом стриме.
Синтаксис
public function groups(): arrayПараметры
Нет.
Возвращает
Массив групп, как их видит сервер: name, consumers, pending, last-delivered-id,
lag.
Пример
foreach ($events->groups() as $group) {
if ($group['pending'] > 1000) {
$this->alert("группа {$group['name']} не успевает");
}
}lag — сколько записей группа ещё не видела; вместе с pending это два разных
отставания: первое про «не прочитано», второе про «прочитано, но не подтверждено».
Ключ целиком
Отдельных сроков жизни у записей нет — рост ограничивают подрезкой. Срок задаётся всему ключу.
expireKey()
Задаёт срок жизни всему журналу.
Синтаксис
public function expireKey(int $seconds): boolПараметры
$seconds — через сколько секунд удалить ключ целиком, считая от текущего момента.
Возвращает
true, если срок установлен; false, если ключа нет.
keyTtl()
Сообщает, сколько журналу осталось жить.
Синтаксис
public function keyTtl(): ?intПараметры
Нет.
Возвращает
Оставшиеся секунды, либо null — и когда срока нет, и когда нет ключа.
deleteKey()
Удаляет журнал вместе с группами.
Синтаксис
public function deleteKey(): boolПараметры
Нет.
Возвращает
true, если ключ существовал.
Уносит и записи, и группы с их учётом незавершённого.
name()
Возвращает имя ключа с префиксом — то, что видит сервер.
Синтаксис
public function name(): stringПараметры
Нет.
Возвращает
Полное имя ключа. Нужно для команд, которых нет в ручке: их выполняют через
raw().
close()
Закрывает соединение, которое открыл follow().
Вызывать безопасно всегда: если блокирующего чтения не было, метод ничего не делает.
Синтаксис
public function close(): voidПараметры
Нет.
Возвращает
Ничего.
Группы потребителей
Без группы каждый читатель видит все записи и сам следит за своей позицией. Группа меняет обе части: запись уходит одному потребителю группы, и сервер держит её в списке незавершённых, пока её не подтвердят.
Отсюда гарантия, ради которой это и берут: воркер взял запись и умер — она не потеряна.
Она числится за ним, её видно в pending(), и другой воркер может её забрать
через claimStale().
$events->ensureGroup('mailers'); // один раз, при старте
$group = $events->group('mailers', consumer: 'worker-1');Создание
ensureGroup()
Создаёт группу потребителей, если её ещё нет.
Стрим при необходимости создаётся тоже, поэтому группу можно объявить на старте приложения, до первой записи.
Синтаксис
public function ensureGroup(string $name, string $from = '0'): boolПараметры
$name — имя группы.
$from — откуда группа начинает читать. По умолчанию '0' — с начала журнала, то есть
группа получит всю накопленную историю. '$' — только записи, добавленные после создания
группы.
Возвращает
true, если группа создана; false, если она уже была — второе не ошибка, и повторный
вызов при каждом старте безопасен.
Пример
$events->ensureGroup('mailers'); // с начала журнала
$events->ensureGroup('metrics', '$'); // только новоеПочему группа не создаётся сама при `consume()`
Тогда опечатка в имени превращалась бы в новую пустую группу, которая молча ничего не
получает. Явное создание оставляет опечатке единственный исход — ошибку NOGROUP,
называющую имя.
group()
Возвращает ручку на группу от лица конкретного потребителя.
Ничего не создаёт и на сервер не ходит — как и остальные ручки, это представление.
Синтаксис
public function group(string $name, string $consumer = 'default'): RedisStreamGroupПараметры
$name — имя группы. Она должна существовать: создаёт её
ensureGroup().
$consumer — имя потребителя, от лица которого будут идти команды. По умолчанию
'default'.
Возвращает
RedisStreamGroup.
Пример
$group = $events->group('mailers', consumer: 'worker-1');Справочник RedisStreamGroup
RedisStreamGroup — ручка на пару «группа + потребитель». Имя потребителя в ней не
случайно: сервер ведёт список незавершённых записей на потребителя, а не на группу,
поэтому все методы восстановления работают от чьего-то лица.
Потребление
consume()
Забирает записи, которые ещё никому в группе не выдавались.
Каждая запись достаётся ровно одному потребителю группы и попадает в его список незавершённых, пока не будет подтверждена.
Синтаксис
public function consume(int $count = 10, float $timeout = 0.0): arrayПараметры
$count — не больше стольких записей за вызов. По умолчанию 10.
$timeout — сколько секунд ждать новых записей. По умолчанию 0.0 — ждать
неограниченно. Отрицательное значение означает «не ждать вовсе»: заглянуть и сразу
вернуться.
Возвращает
Массив StreamEntry; пустой, если новых записей нет.
Ошибки
RedisCommandException с текстом NOGROUP — если группы не существует.
Пример
$group->consume(count: 10, timeout: 5); // ждёт, своё соединение
$group->consume(timeout: -1); // заглянуть и вернуться, соединение из пулаОтрицательный таймаут — единственный режим, который не берёт отдельное соединение; он для проверок и разовых заглядываний.
backlog()
Возвращает то, что этот потребитель уже взял, но не подтвердил.
Это первое, что должен сделать перезапустившийся воркер: записи, взятые до перезапуска,
числятся за ним и через consume() больше не придут. Прочитать их обратно —
способ продолжить работу, не дожидаясь, пока их кто-то перехватит по простою.
Синтаксис
public function backlog(int $count = 10): arrayПараметры
$count — не больше стольких записей. По умолчанию 10.
Возвращает
Массив StreamEntry. Не ждёт: если незавершённого нет, сразу
пустой массив.
Пример
$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()
Отмечает записи обработанными, убирая их из списка незавершённых.
Синтаксис
public function ack(StreamEntry|string ...$entries): intПараметры
...$entries — записи либо их идентификаторы, вперемешку и сколько угодно. Можно не
передать ни одной.
Возвращает
Число записей, которые числились незавершёнными и теперь подтверждены. Без аргументов —
0.
Пример
$group->ack($entry); // объект
$group->ack($entry->id); // или идентификатор
$group->ack(...$entries); // или пачкойПодтверждение не удаляет запись из журнала
ack() закрывает учёт в группе: запись перестаёт числиться незавершённой. Сама запись
остаётся в стриме и по-прежнему видна в range() — на то он и журнал.
Ограничивать рост нужно отдельно, подрезкой.
Подтверждать автоматически, при выдаче, пакет не станет: это уничтожило бы единственную гарантию группы — что запись переживёт воркера, который её взял и не закончил.
Восстановление
pending()
Возвращает записи, выданные группой и никем не подтверждённые.
Синтаксис
public function pending(int $count = 100, ?string $consumer = null): arrayПараметры
$count — не больше стольких записей. По умолчанию 100.
$consumer — ограничить одним потребителем. По умолчанию null — по всей группе.
Возвращает
Массив PendingEntry — это не сами записи, а их учётные карточки: идентификатор, за кем
числится, сколько простаивает и сколько раз выдавалась.
Пример
foreach ($group->pending() as $stuck) {
$stuck->id; // '1755600000123-0'
$stuck->consumer; // 'worker-1'
$stuck->idleMs; // 61240 — сколько миллисекунд назад её выдали
$stuck->deliveries; // 3 — сколько раз выдавали
}Две цифры отвечают на разные вопросы. Большой idleMs при одной доставке — потребитель
умер, держа запись. Растущий deliveries — запись роняет всех, кто за неё берётся.
Что делать со второй — перезапустить, разбудить алерт, отложить в «мёртвые» — решает приложение: цену повтора знает только автор задачи.
pendingCount()
Сообщает, сколько всего записей группа выдала и не получила подтверждения.
Синтаксис
public function pendingCount(): intПараметры
Нет.
Возвращает
Число незавершённых записей по всей группе.
claimStale()
Забирает себе записи, слишком долго висящие у другого потребителя.
Синтаксис
public function claimStale(int $idle, int $count = 10, string $from = '0-0'): arrayПараметры
$idle — сколько миллисекунд запись должна простаивать, чтобы её можно было забрать.
Порог подбирают по самой долгой честной обработке: если письмо отправляется до тридцати
секунд, idle: 60_000 не отберёт работу у живого воркера.
$count — не больше стольких записей. По умолчанию 10.
$from — с какого места просматривать список незавершённых. По умолчанию '0-0' — с
начала.
Возвращает
Массив StreamEntry — уже забранные записи, готовые к обработке.
Пример
foreach ($group->claimStale(idle: 60_000) as $entry) {
$this->handle($entry);
$group->ack($entry);
}Повторный вызов не вернёт те же записи
Захват сбрасывает счётчик простоя, поэтому следующий вызов с тем же порогом их уже не увидит — цикл сходится, а не крутится на одном месте. Счётчик доставок при этом растёт, и по нему видно, что запись пошла по второму кругу.
consumers()
Возвращает всех потребителей группы, как их видит сервер.
Синтаксис
public function consumers(): arrayПараметры
Нет.
Возвращает
Массив потребителей с полями name, pending, idle.
Пример
foreach ($group->consumers() as $consumer) {
if ($consumer['idle'] > 300_000 && $consumer['pending'] > 0) {
$this->alert("{$consumer['name']} молчит пять минут и держит работу");
}
}destroy()
Удаляет группу вместе с её учётом и потребителями.
Записи журнала при этом остаются — удаляется только учёт того, кто что прочитал.
Синтаксис
public function destroy(): boolПараметры
Нет.
Возвращает
true, если группа существовала.
Служебное
as()
Возвращает ту же группу от лица другого потребителя.
Исходную ручку не меняет — создаёт новую. Нужен там, где один процесс работает за нескольких: например, обходит чужие зависшие записи.
Синтаксис
public function as(string $consumer): RedisStreamGroupПараметры
$consumer — имя потребителя.
Возвращает
Новую RedisStreamGroup с тем же стримом и группой.
Пример
$group->as('worker-2')->backlog(); // что не доделал второй воркерRedisStreamGroup::name()
Возвращает имя группы.
Синтаксис
public function name(): stringВозвращает
Имя группы — то, что передали в group().
consumer()
Возвращает имя потребителя, от лица которого работает ручка.
Синтаксис
public function consumer(): stringВозвращает
Имя потребителя.
Имена потребителей должны быть уникальны
Сервер ведёт список незавершённых на потребителя. Два воркера с одним именем делят один список и будут разбирать работу друг друга — включая ту, что первый прямо сейчас обрабатывает. Называйте по чему-то устойчивому и различимому: слот воркера, хост и pid.
RedisStreamGroup::close()
Закрывает соединение, открытое блокирующим consume().
Синтаксис
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() - Пул соединений — почему блокирующие чтения держатся в стороне от пула