Операции с потоками
В Managed Valkey все управление потоками выполняется командами Valkey. С полным списком команд Valkey вы можете ознакомиться в официальной документации Valkey.
Обратите внимание, что в Managed Valkey некоторые команды Valkey недоступны.
Получить записи по диапазону
Заголовок раздела «Получить записи по диапазону»Чтобы получить записи по диапазону идентификаторов, используйте одну из команд:
- XRANGE — получить записи по возрастанию идентификаторов;
- XREVRANGE — получить записи по убыванию идентификаторов.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните нужную команду:
- XRANGE
- XREVRANGE
text XRANGE <имя ключа> <начало> <конец> [COUNT <количество>]Для чтения от начала до конца потока задайте границы
- +.В командах:
- Границы — идентификаторы вида
1700000000000-0. Обе границы будут включены. Значения-и+обозначают отсутствие нижнего и верхнего ограничения. COUNT <количество>ограничит число записей в ответе положительным целым числом.
Команда вернет идентификаторы записей и их поля со значениями в выбранном порядке. Если ключ не существует или диапазон пуст, вернется пустой список. Записи сохранятся.
Для чтения по частям повторяйте команду с
COUNT. В следующем вызове заменяйте движущуюся границу на последний полученный идентификатор с префиксом(, например(1700000000000-0, чтобы исключить повтор этой записи. ДляXRANGEменяйте начало, дляXREVRANGE— конец.
Получить количество записей
Заголовок раздела «Получить количество записей»Чтобы получить количество записей в потоке, используйте команду XLEN.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните команду:
text XLEN <имя ключа>Команда вернет количество записей, которые сейчас хранятся в потоке. Если ключ не существует или поток пуст, вернется
0. Этот ответ не будет количеством необработанных записей: обработанные записи могут оставаться в потоке.Пустой поток может существовать как отдельный ключ. Для проверки наличия ключа используйте EXISTS.
Получить информацию о потоке
Заголовок раздела «Получить информацию о потоке»Чтобы получить информацию о потоке, используйте команду XINFO STREAM.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните команду:
text XINFO STREAM <имя ключа>В ответе будут выведены сведения о потоке, в том числе:
length— текущее количество записей;groups— количество групп потребителей;last-generated-id— последний сгенерированный идентификатор;entries-added— счетчик добавленных записей, включая впоследствии удаленные;first-entryиlast-entry— первая и последняя сохранившиеся записи.
Если поток пуст, вместо первой и последней записей вернется
nil. Если ключ не существует, команда вернет ошибку. Для получения списка групп используйте XINFO GROUPS.
Добавить запись
Заголовок раздела «Добавить запись»Чтобы добавить запись в поток, используйте команду XADD.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните команду:
text XADD <имя ключа> [NOMKSTREAM] <идентификатор> <поле 1> <значение 1> [<поле 2> <значение 2> ...]В команде:
<идентификатор>—*для автоматического формирования или явно заданное значение вида1700000000000-0.<поле>и<значение>— данные записи. Укажите хотя бы одну пару.NOMKSTREAM— не создавать поток, если ключ отсутствует.
Автоматический идентификатор будет сформирован с учетом времени сервера и порядкового номера. Для явно заданного идентификатора используйте значение больше
0-0и большеlast-generated-idиз информации о потоке, даже если запись с прежним идентификатором уже удалена.Команда добавит запись и вернет ее идентификатор. Если ключ не существует, будет создан поток. При
NOMKSTREAMи отсутствующем ключе вернетсяnil. Недопустимый идентификатор приведет к ошибке.Для удаления старых записей используйте инструкцию по обрезке потока.
Читать записи после указанного идентификатора
Заголовок раздела «Читать записи после указанного идентификатора»Чтобы читать записи после заданной позиции и при необходимости ожидать новые, используйте команду XREAD.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имена исходных ключей с помощью итеративного обхода.
Проверьте тип значений исходных ключей. Для каждого существующего ключа должно вернуться значение
stream.Выберите начальную позицию для каждого потока. Для чтения с начала используйте
0-0, для продолжения — последний полученный идентификатор.Выполните команду:
text XREAD [COUNT <количество>] [BLOCK <таймаут>] STREAMS <имя ключа 1> [<имя ключа 2> ...] <идентификатор 1> [<идентификатор 2> ...]В команде сначала перечислите все ключи, затем столько же идентификаторов в том же порядке.
COUNT— положительное максимальное число записей от каждого потока за вызов.BLOCK— неотрицательное время ожидания в миллисекундах.0задаст ожидание без ограничения времени со стороны команды.
Команда вернет записи с идентификаторами больше указанных, сгруппированные по потокам. Если подходящих записей нет и не указан
BLOCK, сразу вернетсяnil. Если подходящих записей нет и указанBLOCK, соединение будет ожидать данные до таймаута. Ожидание не заблокирует другие соединения. Чтение не удалит записи и не добавит их в PEL.Продолжайте чтение, подставляя для каждого потока последний полученный идентификатор. Если от потока ничего не пришло, сохраните его прежнюю позицию.
Чтобы при первом вызове пропустить имеющиеся записи и ожидать только новые, используйте
$вместе сBLOCK. После получения данных замените$на последний полученный идентификатор: повторное использование$может пропустить записи между вызовами.
Получить список групп потребителей
Заголовок раздела «Получить список групп потребителей»Чтобы получить список групп потребителей потока, используйте команду XINFO GROUPS.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните команду:
text XINFO GROUPS <имя ключа>Для каждой группы будут выведены параметры:
name— имя группы.consumers— количество потребителей.pending— количество записей в PEL.lag— оценка количества еще не выданных группе записей. Не включает записи PEL и может бытьnil, если оценку нельзя вычислить.last-delivered-id— последняя позиция выдачи, и другие сведения.
Для поиска уже доставленных, но не подтвержденных записей используйте XPENDING.
Если групп нет, вернется пустой список. Если ключ не существует, команда вернет ошибку.
Создать группу потребителей
Заголовок раздела «Создать группу потребителей»Чтобы создать группу потребителей, используйте команду XGROUP CREATE.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выберите имя новой группы. Для существующего потока проверьте, что группы с таким именем еще нет.
Выполните команду:
text XGROUP CREATE <имя ключа> <имя группы> <идентификатор> [MKSTREAM]В команде:
<идентификатор>— начальная позиция выдачи. При чтении новых для группы записей будут выдаваться идентификаторы строго больше нее.0-0— начать со всех сохранившихся записей потока.$— пропустить имеющиеся записи и начать с добавленных после создания группы.MKSTREAM— создать пустой поток, если ключ не существует.
Команда создаст группу и вернет
OK. Если группа с таким именем уже существует, вернется ошибкаBUSYGROUP. Если ключ отсутствует иMKSTREAMне указан, команда вернет ошибку. Существующие записи потока сохранятся.
Получить список потребителей группы
Заголовок раздела «Получить список потребителей группы»Чтобы получить список потребителей группы, используйте команду XINFO CONSUMERS.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы.
Выполните команду:
text XINFO CONSUMERS <имя ключа> <имя группы>Для каждого потребителя будут выведены:
name— имя;pending— количество доставленных ему, но не подтвержденных записей;idle— время с последней попытки взаимодействия в миллисекундах;inactive— время с последнего успешного взаимодействия в миллисекундах.
Если потребителей нет, вернется пустой список. Если поток или группа не существует, команда вернет ошибку. Для получения идентификаторов записей PEL используйте XPENDING.
Создать потребителя группы
Заголовок раздела «Создать потребителя группы»Чтобы явно создать потребителя группы, используйте команду XGROUP CREATECONSUMER.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выберите имя потребителя. Имена потребителей одной группы должны различаться, если за ними стоят разные обработчики.
Выполните команду:
text XGROUP CREATECONSUMER <имя ключа> <имя группы> <имя потребителя>Команда вернет
1, если потребитель будет создан, или0, если он уже существует. Потребителю не будут выданы записи. Если поток или группа не существует, команда вернет ошибку.Явное создание не обязательно. При выдаче данных через XREADGROUP потребитель также будет создан автоматически.
Читать записи в составе группы
Заголовок раздела «Читать записи в составе группы»Чтобы получать записи от имени потребителя группы, используйте команду XREADGROUP.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Выполните команду:
text XREADGROUP GROUP <имя группы> <имя потребителя> [COUNT <количество>] [BLOCK <таймаут>] [NOACK] STREAMS <имя ключа> <идентификатор>В команде:
>— получить записи, которые еще не выдавались этой группе.0-0или другой числовой идентификатор — получить собственные неподтвержденные записи с большими идентификаторами.BLOCKиNOACKв этом режиме будут проигнорированы.COUNT— положительное максимальное число записей за вызов.BLOCK— таймаут в миллисекундах;0задаст неограниченное ожидание.NOACK— не добавлять новые записи в PEL. Используйте только при допустимости потери обработки после выдачи.
Команда вернет записи. Новые записи без
NOACKбудут закреплены за потребителем в PEL. Если новых записей для режима>нет, вернетсяnilсразу или после таймаута. При чтении собственной истории без результатов вернется пустой список записей. Если поток или группа отсутствует, вернется ошибкаNOGROUP. При чтении PEL вместо содержимого удаленной записи вернетсяnil.Обработайте полученные записи. Только после успешной обработки подтвердите их.
При восстановлении обработчика сначала прочитайте его PEL с
0-0, продолжая с последнего полученного идентификатора, затем перейдите к>. Записи другого потребителя можно передать или перераспределить. Повторная доставка возможна. Обработка одной записи несколько раз не должна повторно выполнять нежелательное действие.
Получить список неподтвержденных записей
Заголовок раздела «Получить список неподтвержденных записей»Чтобы получить сведения о записях PEL, используйте команду XPENDING.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы.
Выполните нужную команду:
- Общая информация
- Список записей
text XPENDING <имя ключа> <имя группы>Команда вернет количество записей PEL, минимальный и максимальный идентификаторы, а также потребителей с количеством закрепленных записей.
Состояние PEL не изменится. Для чтения содержимого используйте XRANGE, задав идентификатор обеими границами. Для постраничного просмотра повторяйте расширенную форму с началом
(<последний идентификатор>. При отсутствии записей расширенная форма вернет пустой список. Отсутствующая группа или поток приведет к ошибке.
Подтвердить обработку записей
Заголовок раздела «Подтвердить обработку записей»Чтобы подтвердить успешную обработку записей группы, используйте команду XACK.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы.
Возьмите идентификаторы успешно обработанных записей из ответа XREADGROUP, XCLAIM или XAUTOCLAIM. Состав PEL можно проверить.
Выполните команду:
text XACK <имя ключа> <имя группы> <идентификатор 1> [<идентификатор 2> ...]Команда удалит указанные идентификаторы из PEL этой группы и вернет количество подтвержденных записей. Идентификаторы, которых уже нет в PEL, не будут учтены в результате.
Записи потока сохранятся. Подтверждение одной группы не подтвердит обработку в других группах. Не подтверждайте запись только на основании ее получения. Подтверждение должно следовать за успешной обработкой.
Передать неподтвержденные записи другому потребителю
Заголовок раздела «Передать неподтвержденные записи другому потребителю»Чтобы передать выбранные записи PEL другому потребителю, используйте команду XCLAIM.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы.
Получите идентификаторы записей, требующих повторной обработки.
Выполните команду:
text XCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> <идентификатор 1> [<идентификатор 2> ...] [JUSTID]В команде
<минимальное время>— неотрицательный порог времени с последней доставки в миллисекундах. Выберите его так, чтобы не перехватывать записи у нормально работающего обработчика.JUSTIDограничит ответ идентификаторами.Команда передаст подходящие записи указанному потребителю, сбросит их время с последней доставки и вернет переданные записи. Без
JUSTIDсчетчик доставок будет увеличен. Записи вне PEL или не достигшие порога в ответ не попадут. Для записей, уже удаленных из потока, соответствующие ссылки будут удалены из PEL.После успешной обработки подтвердите записи. Передача самостоятельно не подтвердит обработку и не удалит записи из потока.
Итеративно перераспределить неподтвержденные записи
Заголовок раздела «Итеративно перераспределить неподтвержденные записи»Чтобы найти и передать записи PEL по времени с последней доставки, используйте команду XAUTOCLAIM.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Начните обход с идентификатора
0-0:text XAUTOCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> 0-0 [COUNT <количество>] [JUSTID]В команде:
<минимальное время>— неотрицательный порог с последней доставки в миллисекундах.COUNT— задать положительное максимальное число передаваемых записей. По умолчанию —100.JUSTID— ограничить данные переданных записей их идентификаторами.
Вернутся позиция для продолжения обхода, переданные записи и идентификаторы удаленных из потока записей, ссылки на которые были очищены из PEL. У переданных записей время с последней доставки будет сброшено. Если
JUSTIDне указан, счетчик доставок будет увеличен.Повторяйте вызов с полученной позицией вместо
0-0, сохраняя ключ, группу, потребителя и порог времени:text XAUTOCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> <следующий идентификатор> [COUNT <количество>] [JUSTID]Текущий обход завершится при возвращенной позиции
0-0. Пустой список переданных записей при другой позиции не означает завершение. Нулевой курсор не означает, что PEL пуст. Неподходящие по времени записи могут стать доступными для передачи позже.После успешной обработки подтвердите переданные записи.
Изменить позицию выдачи группы
Заголовок раздела «Изменить позицию выдачи группы»Чтобы изменить позицию выдачи новых для группы записей, используйте команду XGROUP SETID.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы и ее текущую позицию
last-delivered-id.Выберите новую позицию. При необходимости получите идентификатор записи.
Выполните команду:
text XGROUP SETID <имя ключа> <имя группы> <идентификатор>В команде
<идентификатор>— новая последняя позиция выдачи. Значение0-0вернет позицию к началу, а$переместит ее к последней записи потока.Команда изменит позицию и вернет
OK. Последующие вызовы XREADGROUP с>будут получать записи с большими идентификаторами. Содержимое потока и существующий PEL сохранятся. Если поток или группа отсутствует, команда вернет ошибку.
Удалить потребителя группы
Заголовок раздела «Удалить потребителя группы»Чтобы удалить неиспользуемого потребителя группы, используйте команду XGROUP DELCONSUMER.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Проверьте его показатель
pending. Если он не равен0, передайте неподтвержденные записи другому потребителю или подтвердите уже успешно обработанные.Убедитесь, что приложение больше не использует имя удаляемого потребителя и его
pendingравен0.Выполните команду:
text XGROUP DELCONSUMER <имя ключа> <имя группы> <имя потребителя>Команда удалит потребителя и вернет число записей, которые числились за ним в PEL перед удалением. Если неподтвержденных записей нет, вернется
0. Если потребитель уже отсутствует, также вернется0. Если поток или группа не существует, команда вернет ошибку.Записи потока сохранятся. При последующей выдаче данных через
XREADGROUPс тем же именем потребитель сможет появиться снова.
Удалить группу потребителей
Заголовок раздела «Удалить группу потребителей»Чтобы удалить группу потребителей, используйте команду XGROUP DESTROY.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте имя группы.
Проверьте неподтвержденные записи и потребителей. Выполняйте удаление, когда приложения больше не используют эту группу.
Выполните команду:
text XGROUP DESTROY <имя ключа> <имя группы>Команда удалит группу вместе с ее потребителями и PEL и вернет
1. Если группы нет, вернется0. Ключ потока, его записи и другие группы сохранятся.Повторное создание группы с тем же именем не восстановит прежний PEL и состояние потребителей.
Обрезать поток
Заголовок раздела «Обрезать поток»Чтобы удалить старые записи по количеству или идентификатору, используйте команду XTRIM.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Просмотрите записи и выберите границу хранения. Если используются группы, проверьте их состояние и неподтвержденные записи.
Выполните нужную команду:
- По количеству
- По идентификатору
text XTRIM <имя ключа> MAXLEN [= | ~] <количество>В команде
<количество>— неотрицательное число записей, которые нужно оставить. При точной обрезке будут сохранены не более этого количества самых новых записей.Без оператора или с
=обрезка будет точной. С~будут удаляться целые внутренние блоки, поэтому часть лишних записей может остаться. Только для режима~можно добавитьLIMIT <лимит>— ограничение объема удаления за вызов.0снимет это ограничение.Команда вернет количество удаленных записей. Если удалять нечего или ключ отсутствует, вернется
0. При удалении всех записей сам поток сохранится.
Удалить записи по идентификаторам
Заголовок раздела «Удалить записи по идентификаторам»Чтобы удалить отдельные записи потока, используйте команду XDEL.
Подключитесь к кластеру.
Переключитесь на нужную базу данных.
Узнайте имя ключа с помощью итеративного обхода.
Проверьте тип значения ключа. Для существующего ключа должно вернуться значение
stream.Узнайте идентификаторы удаляемых записей.
Выполните команду:
text XDEL <имя ключа> <идентификатор 1> [<идентификатор 2> ...]Команда удалит указанные записи и вернет количество удаленных записей. Несуществующие идентификаторы будут пропущены. Если ключ не существует, вернется
0. При удалении последней записи ключ потока сохранится.Удаление не заменит подтверждение обработки. Если запись еще учтена в PEL, ее идентификатор может остаться там без содержимого.