Стримы
Стрим — журнал с добавлением в конец: у каждой записи есть идентификатор, присвоенный сервером, и набор полей. От списка он отличается тем, что чтение не забирает запись, а группы потребителей заставляют сам Redis вести учёт: кто что взял и что ещё не подтверждено.
Что такое стрим
Запись выглядит так: идентификатор 1755600000123-0 (миллисекунды серверных часов и
порядковый номер внутри этой миллисекунды) и поля {"type": "signup", "user": "42"}.
Идентификаторы монотонно растут, поэтому «прочитать всё после вот этого» — обычная
операция, а не поиск.
Свойства, из которых следует всё остальное:
- Чтение ничего не изымает. Запись живёт, пока её не удалят или не подрежут журнал. Одну и ту же запись могут прочитать сколько угодно независимых читателей.
- Порядок гарантирован и совпадает с порядком добавления.
- Журнал растёт вечно, если его не ограничить. Это главная плата за историю.
- Группы ведёт сервер. Он помнит, какая запись какому потребителю ушла и была ли подтверждена. Список этого не умеет — там учёт приходится строить руками.
Чем это отличается от списка
| Список | Стрим | |
|---|---|---|
| После чтения | элемент исчез | запись на месте |
| Читателей на сообщение | один | сколько угодно |
| История | нет | есть, пока не подрезали |
| Учёт «взято, но не завершено» | вручную (moveTo + remove) |
встроенный, в группе |
| Рост | ограничен потреблением | вечный, пока не ограничить |
| Идентификаторы | нет | есть, по ним читают и подтверждают |
Если нужна очередь и история не нужна — берите список: он проще и дешевле. Стрим берут ради истории, нескольких независимых читателей или встроенного учёта доставки.
Где принято применять
Шина событий. Одно событие — несколько независимых обработчиков: письмо, метрика, вебхук. Каждый читает своей группой и не мешает остальным.
Очередь с гарантией. Воркер взял запись, упал — она осталась в списке
незавершённых и достанется другому. Ровно то, что в списках приходится собирать из
moveTo() и remove().
Журнал изменений. Аудит, история заказа, лента активности: можно перечитать, отмотать назад, разобрать инцидент по следам.
Буфер с догоняющим чтением. Обработчик пишет, потребитель читает в своём темпе и после перезапуска продолжает с того места, где остановился.
Где стрим — неправильный выбор:
| Задача | Почему не стрим | Что вместо |
|---|---|---|
| Простая очередь без истории | лишний рост и лишние понятия | список |
| Уведомить всех подключённых сейчас | стрим хранит, а не рассылает | pub/sub |
| Найти запись по содержимому | поиска нет, только по идентификатору и диапазону | индекс на множествах |
| Хранить состояние объекта | стрим про события, а не про текущее значение | хеш |
Ручка
$events = $store->stream('events'); // сервер увидит 'session:events'stream() ничего не выполняет — это представление, а не запрос. Префикс применяется
здесь один раз. Соединения ручка не держит, кроме одного случая: follow() открывает
собственное, потому что блокирующее чтение занимает соединение на всё время ожидания.
Записи
Методы чтения возвращают 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 — миллисекунды серверных часов
}Драйвер отдаёт записи как ['1755600000123-0' => ['type' => 'signup']] — массив, у
которого единственный ключ и есть идентификатор. Обходить такое приходится через
key()/current(), а идентификатор нужен постоянно: им подтверждают обработку и
запоминают позицию.
Справочник RedisStream
Запись
add()
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,
если ключ занят структурой другого типа.
$events->add(['type' => 'signup', 'user' => '42']);
// '1755600000123-0'
$events->add(['level' => 'warn', 'msg' => $text], cap: 10_000);Свои идентификаторы почти никогда не нужны: серверные монотонны и включают время. Ручное назначение имеет смысл только при переносе данных из другого журнала.
Значения полей — строки
Как и везде в Redis. Массив в поле превратится в "Array", если у конфигурации не задан
сериализатор.
trim()
public function trim(int $cap, bool $exact = false): intОставляет примерно $cap новейших записей.
| Аргумент | Тип | По умолчанию | Что делает |
|---|---|---|---|
$cap |
int |
— | сколько записей оставить |
$exact |
bool |
false |
подрезать ровно, а не приблизительно |
Возвращает число удалённых записей.
Приблизительная подрезка может не удалить ничего
Она останавливается на границе внутреннего блока стрима — так дешевле, и Redis
рекомендует именно её. Следствие неочевидное: на коротком журнале trim(10) вполне может
вернуть 0 и оставить всё как было. Это не сбой; если нужна предсказуемая длина, просите
exact: true и платите за обход.
$events->trim(10_000); // дёшево, длина «примерно такая»
$events->trim(10_000, exact: true); // ровно столько, дорожеtrimBefore()
public function trimBefore(string $id, bool $exact = false): intУдаляет всё, что старше указанного идентификатора.
| Аргумент | Тип | По умолчанию | Что делает |
|---|---|---|---|
$id |
string |
— | граница; записи с меньшим идентификатором удаляются |
$exact |
bool |
false |
точная граница вместо приблизительной |
Возвращает число удалённых записей.
Так журнал ограничивают по времени, а не по количеству: идентификатор начинается с миллисекунд серверных часов, поэтому «старше суток» — это идентификатор:
$events->trimBefore((string) ((time() - 86400) * 1000));delete()
public function delete(string ...$ids): intУдаляет записи по идентификаторам. Возвращает число существовавших. Без аргументов —
0.
Удаление из середины журнала допустимо, но это исключение: обычный способ ограничить
стрим — подрезка. delete() нужен, когда конкретную запись нельзя больше показывать.
Чтение
count()
public function count(): intЧисло записей; для несуществующего ключа — 0. Аргументов нет.
range()
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>; для пустого журнала — пустой массив.
$events->range(); // всё
$events->range(count: 100); // первая сотня
$events->range($since, count: 50); // с известной позицииreverse()
public function reverse(?int $count = null): arrayТо же самое, но новейшие первыми — дешёвый способ посмотреть, что только что произошло.
$events->reverse(count: 10); // последние десять, свежие сверхуafter()
public function after(string $id, int $count = 10): arrayЗаписи, добавленные после указанного идентификатора, без ожидания.
| Аргумент | Тип | По умолчанию | Что делает |
|---|---|---|---|
$id |
string |
— | позиция, после которой читать |
$count |
int |
10 |
не больше стольких записей |
Возвращает list<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 |
float |
0.0 |
сколько секунд ждать; 0 — неограниченно |
$count |
int |
10 |
не больше стольких записей за раз |
$from |
?string |
null |
позиция для первого вызова; '0' — начать со всей истории |
Возвращает list<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 это два разных
отставания: первое про «не прочитано», второе про «прочитано, но не подтверждено».
Ключ целиком
$events->expireKey(86400); // срок жизни всему журналу
$events->keyTtl(); // 86400 | null
$events->deleteKey(); // удалить журнал вместе с группами
$events->name(); // 'session:events'
$events->close(); // закрыть соединение, открытое follow()deleteKey() уносит и записи, и группы с их учётом. Отдельных сроков жизни у записей
нет — ограничивают подрезкой.
Группы потребителей
Зачем они
Без группы каждый читатель видит все записи и сам следит за своей позицией. Группа меняет обе части: запись уходит одному потребителю группы, и сервер держит её в списке незавершённых, пока её не подтвердят.
Отсюда гарантия, ради которой это и берут: воркер взял запись и умер — она не потеряна.
Она числится за ним, её видно в pending(), и другой воркер может её забрать.
$events->ensureGroup('mailers'); // один раз, при старте
$group = $events->group('mailers', consumer: 'worker-1');ensureGroup()
public function ensureGroup(string $name, string $from = '0'): bool| Аргумент | Тип | По умолчанию | Что делает |
|---|---|---|---|
$name |
string |
— | имя группы |
$from |
string |
'0' |
откуда группа начинает читать: '0' — с начала журнала, '$' — только новое |
Возвращает true, если группа создана, и false, если она уже была — второе не
ошибка. Стрим при необходимости создаётся тоже, поэтому группу можно объявить на старте,
до первой записи.
Почему группа не создаётся сама при consume()
Тогда опечатка в имени превращалась бы в новую пустую группу, которая молча ничего не
получает. Явное создание оставляет опечатке единственный исход — ошибку NOGROUP,
называющую имя.
group()
public function group(string $name, string $consumer = 'default'): RedisStreamGroupРучка на группу от лица конкретного потребителя. Ничего не создаёт и не отправляет.
Справочник RedisStreamGroup
Потребление
consume()
public function consume(int $count = 10, float $timeout = 0.0): arrayЗабирает записи, которые ещё никому в группе не выдавались.
| Аргумент | Тип | По умолчанию | Что делает |
|---|---|---|---|
$count |
int |
10 |
не больше стольких записей |
$timeout |
float |
0.0 |
секунд ожидания; 0 — неограниченно, отрицательное — не ждать вовсе |
Возвращает list<StreamEntry>. Бросает RedisCommandException с NOGROUP,
если группы не существует.
Каждая запись достаётся ровно одному потребителю группы и попадает в его список незавершённых, пока не будет подтверждена.
$group->consume(count: 10, timeout: 5); // ждёт, своё соединение
$group->consume(timeout: -1); // заглянуть и вернуться, соединение из пулаОтрицательный таймаут — единственный режим, который не берёт отдельное соединение; он для проверок и разовых заглядываний.
backlog()
public function backlog(int $count = 10): arrayВозвращает то, что этот потребитель уже взял, но не подтвердил.
Возвращает list<StreamEntry>, не ждёт.
Это первое, что должен сделать перезапустившийся воркер: записи, взятые до перезапуска,
числятся за ним и через consume() больше не придут. Прочитать их обратно — способ
продолжить работу, не дожидаясь, пока их кто-то перехватит по простою.
$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 |
StreamEntry|string |
— | записи или их идентификаторы |
Возвращает число записей, которые числились незавершёнными и теперь подтверждены.
Без аргументов — 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 |
int |
100 |
не больше стольких |
$consumer |
?string |
null |
ограничить одним потребителем; null — по всей группе |
Возвращает list<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 |
int |
— | сколько миллисекунд запись должна простаивать, чтобы её можно было забрать |
$count |
int |
10 |
не больше стольких |
$from |
string |
'0-0' |
с какого места просматривать список незавершённых |
Возвращает list<StreamEntry> — уже забранные записи, готовые к обработке.
foreach ($group->claimStale(idle: 60_000) as $entry) {
$this->handle($entry);
$group->ack($entry);
}Повторный вызов не вернёт те же записи
Захват сбрасывает счётчик простоя, поэтому следующий вызов с тем же порогом их уже не увидит — цикл сходится, а не крутится на одном месте. Счётчик доставок при этом растёт, и по нему видно, что запись пошла по второму кругу.
Порог подбирают по самой долгой честной обработке: если письмо отправляется до тридцати
секунд, idle: 60_000 не отберёт работу у живого воркера.
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Удаляет группу вместе с её учётом и потребителями. Записи журнала остаются.
Служебное
$group->as('worker-2'); // та же группа от лица другого потребителя
$group->name(); // 'mailers'
$group->consumer(); // 'worker-1'
$group->close(); // закрыть соединение, открытое блокирующим consume()Имена потребителей должны быть уникальны
Сервер ведёт список незавершённых на потребителя. Два воркера с одним именем делят один список и будут разбирать работу друг друга — включая ту, что первый прямо сейчас обрабатывает. Называйте по чему-то устойчивому и различимому: слот воркера, хост и pid.
Соответствие командам 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().
Дальше
- Списки — очередь, когда история не нужна
- Хеши — состояние объекта рядом с журналом событий
- Пул соединений — почему блокирующие чтения держатся в стороне от пула