Перейти к содержимому

Операции с потоками

В Managed Valkey все управление потоками выполняется командами Valkey. С полным списком команд Valkey вы можете ознакомиться в официальной документации Valkey.

Обратите внимание, что в Managed Valkey некоторые команды Valkey недоступны.

Чтобы получить записи по диапазону идентификаторов, используйте одну из команд:

  • XRANGE — получить записи по возрастанию идентификаторов;
  • XREVRANGE — получить записи по убыванию идентификаторов.
  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Выполните нужную команду:

    • XRANGE
    • XREVRANGE
    text
    XRANGE <имя ключа> <начало> <конец> [COUNT <количество>]

    Для чтения от начала до конца потока задайте границы - +.

    В командах:

    • Границы — идентификаторы вида 1700000000000-0. Обе границы будут включены. Значения - и + обозначают отсутствие нижнего и верхнего ограничения.
    • COUNT <количество> ограничит число записей в ответе положительным целым числом.

    Команда вернет идентификаторы записей и их поля со значениями в выбранном порядке. Если ключ не существует или диапазон пуст, вернется пустой список. Записи сохранятся.

    Для чтения по частям повторяйте команду с COUNT. В следующем вызове заменяйте движущуюся границу на последний полученный идентификатор с префиксом (, например (1700000000000-0, чтобы исключить повтор этой записи. Для XRANGE меняйте начало, для XREVRANGE — конец.

Чтобы получить количество записей в потоке, используйте команду XLEN.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Выполните команду:

    text
    XLEN <имя ключа>

    Команда вернет количество записей, которые сейчас хранятся в потоке. Если ключ не существует или поток пуст, вернется 0. Этот ответ не будет количеством необработанных записей: обработанные записи могут оставаться в потоке.

    Пустой поток может существовать как отдельный ключ. Для проверки наличия ключа используйте EXISTS.

Чтобы получить информацию о потоке, используйте команду XINFO STREAM.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Выполните команду:

    text
    XINFO STREAM <имя ключа>

    В ответе будут выведены сведения о потоке, в том числе:

    • length — текущее количество записей;
    • groups — количество групп потребителей;
    • last-generated-id — последний сгенерированный идентификатор;
    • entries-added — счетчик добавленных записей, включая впоследствии удаленные;
    • first-entry и last-entry — первая и последняя сохранившиеся записи.

    Если поток пуст, вместо первой и последней записей вернется nil. Если ключ не существует, команда вернет ошибку. Для получения списка групп используйте XINFO GROUPS.

Чтобы добавить запись в поток, используйте команду XADD.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

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

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Выполните команду:

    text
    XADD <имя ключа> [NOMKSTREAM] <идентификатор> <поле 1> <значение 1> [<поле 2> <значение 2> ...]

    В команде:

    • <идентификатор> — * для автоматического формирования или явно заданное значение вида 1700000000000-0.
    • <поле> и <значение> — данные записи. Укажите хотя бы одну пару.
    • NOMKSTREAM — не создавать поток, если ключ отсутствует.

    Автоматический идентификатор будет сформирован с учетом времени сервера и порядкового номера. Для явно заданного идентификатора используйте значение больше 0-0 и больше last-generated-id из информации о потоке, даже если запись с прежним идентификатором уже удалена.

    Команда добавит запись и вернет ее идентификатор. Если ключ не существует, будет создан поток. При NOMKSTREAM и отсутствующем ключе вернется nil. Недопустимый идентификатор приведет к ошибке.

    Для удаления старых записей используйте инструкцию по обрезке потока.

Читать записи после указанного идентификатора

Заголовок раздела «Читать записи после указанного идентификатора»

Чтобы читать записи после заданной позиции и при необходимости ожидать новые, используйте команду XREAD.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имена исходных ключей с помощью итеративного обхода.

  4. Проверьте тип значений исходных ключей. Для каждого существующего ключа должно вернуться значение stream.

  5. Выберите начальную позицию для каждого потока. Для чтения с начала используйте 0-0, для продолжения — последний полученный идентификатор.

  6. Выполните команду:

    text
    XREAD [COUNT <количество>] [BLOCK <таймаут>] STREAMS <имя ключа 1> [<имя ключа 2> ...] <идентификатор 1> [<идентификатор 2> ...]

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

    • COUNT — положительное максимальное число записей от каждого потока за вызов.
    • BLOCK — неотрицательное время ожидания в миллисекундах. 0 задаст ожидание без ограничения времени со стороны команды.

    Команда вернет записи с идентификаторами больше указанных, сгруппированные по потокам. Если подходящих записей нет и не указан BLOCK, сразу вернется nil. Если подходящих записей нет и указан BLOCK, соединение будет ожидать данные до таймаута. Ожидание не заблокирует другие соединения. Чтение не удалит записи и не добавит их в PEL.

  7. Продолжайте чтение, подставляя для каждого потока последний полученный идентификатор. Если от потока ничего не пришло, сохраните его прежнюю позицию.

    Чтобы при первом вызове пропустить имеющиеся записи и ожидать только новые, используйте $ вместе с BLOCK. После получения данных замените $ на последний полученный идентификатор: повторное использование $ может пропустить записи между вызовами.

Чтобы получить список групп потребителей потока, используйте команду XINFO GROUPS.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Выполните команду:

    text
    XINFO GROUPS <имя ключа>

    Для каждой группы будут выведены параметры:

    • name — имя группы.
    • consumers — количество потребителей.
    • pending — количество записей в PEL.
    • lag — оценка количества еще не выданных группе записей. Не включает записи PEL и может быть nil, если оценку нельзя вычислить.
    • last-delivered-id — последняя позиция выдачи, и другие сведения.

    Для поиска уже доставленных, но не подтвержденных записей используйте XPENDING.

    Если групп нет, вернется пустой список. Если ключ не существует, команда вернет ошибку.

Чтобы создать группу потребителей, используйте команду XGROUP CREATE.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

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

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

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

  6. Выполните команду:

    text
    XGROUP CREATE <имя ключа> <имя группы> <идентификатор> [MKSTREAM]

    В команде:

    • <идентификатор> — начальная позиция выдачи. При чтении новых для группы записей будут выдаваться идентификаторы строго больше нее.
    • 0-0 — начать со всех сохранившихся записей потока.
    • $ — пропустить имеющиеся записи и начать с добавленных после создания группы.
    • MKSTREAM — создать пустой поток, если ключ не существует.

    Команда создаст группу и вернет OK. Если группа с таким именем уже существует, вернется ошибка BUSYGROUP. Если ключ отсутствует и MKSTREAM не указан, команда вернет ошибку. Существующие записи потока сохранятся.

Чтобы получить список потребителей группы, используйте команду XINFO CONSUMERS.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы.

  6. Выполните команду:

    text
    XINFO CONSUMERS <имя ключа> <имя группы>

    Для каждого потребителя будут выведены:

    • name — имя;
    • pending — количество доставленных ему, но не подтвержденных записей;
    • idle — время с последней попытки взаимодействия в миллисекундах;
    • inactive — время с последнего успешного взаимодействия в миллисекундах.

    Если потребителей нет, вернется пустой список. Если поток или группа не существует, команда вернет ошибку. Для получения идентификаторов записей PEL используйте XPENDING.

Чтобы явно создать потребителя группы, используйте команду XGROUP CREATECONSUMER.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы. Если группы нет, создайте ее.

  6. Выберите имя потребителя. Имена потребителей одной группы должны различаться, если за ними стоят разные обработчики.

  7. Выполните команду:

    text
    XGROUP CREATECONSUMER <имя ключа> <имя группы> <имя потребителя>

    Команда вернет 1, если потребитель будет создан, или 0, если он уже существует. Потребителю не будут выданы записи. Если поток или группа не существует, команда вернет ошибку.

    Явное создание не обязательно. При выдаче данных через XREADGROUP потребитель также будет создан автоматически.

Чтобы получать записи от имени потребителя группы, используйте команду XREADGROUP.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы. Если группы нет, создайте ее.

  6. Узнайте имя потребителя или создайте нового.

  7. Выполните команду:

    text
    XREADGROUP GROUP <имя группы> <имя потребителя> [COUNT <количество>] [BLOCK <таймаут>] [NOACK] STREAMS <имя ключа> <идентификатор>

    В команде:

    • > — получить записи, которые еще не выдавались этой группе.
    • 0-0 или другой числовой идентификатор — получить собственные неподтвержденные записи с большими идентификаторами. BLOCK и NOACK в этом режиме будут проигнорированы.
    • COUNT — положительное максимальное число записей за вызов.
    • BLOCK — таймаут в миллисекундах; 0 задаст неограниченное ожидание.
    • NOACK — не добавлять новые записи в PEL. Используйте только при допустимости потери обработки после выдачи.

    Команда вернет записи. Новые записи без NOACK будут закреплены за потребителем в PEL. Если новых записей для режима > нет, вернется nil сразу или после таймаута. При чтении собственной истории без результатов вернется пустой список записей. Если поток или группа отсутствует, вернется ошибка NOGROUP. При чтении PEL вместо содержимого удаленной записи вернется nil.

  8. Обработайте полученные записи. Только после успешной обработки подтвердите их.

    При восстановлении обработчика сначала прочитайте его PEL с 0-0, продолжая с последнего полученного идентификатора, затем перейдите к >. Записи другого потребителя можно передать или перераспределить. Повторная доставка возможна. Обработка одной записи несколько раз не должна повторно выполнять нежелательное действие.

Получить список неподтвержденных записей

Заголовок раздела «Получить список неподтвержденных записей»

Чтобы получить сведения о записях PEL, используйте команду XPENDING.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы.

  6. Выполните нужную команду:

    • Общая информация
    • Список записей
    text
    XPENDING <имя ключа> <имя группы>

    Команда вернет количество записей PEL, минимальный и максимальный идентификаторы, а также потребителей с количеством закрепленных записей.

    Состояние PEL не изменится. Для чтения содержимого используйте XRANGE, задав идентификатор обеими границами. Для постраничного просмотра повторяйте расширенную форму с началом (<последний идентификатор>. При отсутствии записей расширенная форма вернет пустой список. Отсутствующая группа или поток приведет к ошибке.

Чтобы подтвердить успешную обработку записей группы, используйте команду XACK.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы.

  6. Возьмите идентификаторы успешно обработанных записей из ответа XREADGROUP, XCLAIM или XAUTOCLAIM. Состав PEL можно проверить.

  7. Выполните команду:

    text
    XACK <имя ключа> <имя группы> <идентификатор 1> [<идентификатор 2> ...]

    Команда удалит указанные идентификаторы из PEL этой группы и вернет количество подтвержденных записей. Идентификаторы, которых уже нет в PEL, не будут учтены в результате.

    Записи потока сохранятся. Подтверждение одной группы не подтвердит обработку в других группах. Не подтверждайте запись только на основании ее получения. Подтверждение должно следовать за успешной обработкой.

Передать неподтвержденные записи другому потребителю

Заголовок раздела «Передать неподтвержденные записи другому потребителю»

Чтобы передать выбранные записи PEL другому потребителю, используйте команду XCLAIM.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы.

  6. Получите идентификаторы записей, требующих повторной обработки.

  7. Узнайте имя нового потребителя или создайте его.

  8. Выполните команду:

    text
    XCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> <идентификатор 1> [<идентификатор 2> ...] [JUSTID]

    В команде <минимальное время> — неотрицательный порог времени с последней доставки в миллисекундах. Выберите его так, чтобы не перехватывать записи у нормально работающего обработчика. JUSTID ограничит ответ идентификаторами.

    Команда передаст подходящие записи указанному потребителю, сбросит их время с последней доставки и вернет переданные записи. Без JUSTID счетчик доставок будет увеличен. Записи вне PEL или не достигшие порога в ответ не попадут. Для записей, уже удаленных из потока, соответствующие ссылки будут удалены из PEL.

    После успешной обработки подтвердите записи. Передача самостоятельно не подтвердит обработку и не удалит записи из потока.

Итеративно перераспределить неподтвержденные записи

Заголовок раздела «Итеративно перераспределить неподтвержденные записи»

Чтобы найти и передать записи PEL по времени с последней доставки, используйте команду XAUTOCLAIM.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы и выберите потребителя. При необходимости создайте потребителя.

  6. Начните обход с идентификатора 0-0:

    text
    XAUTOCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> 0-0 [COUNT <количество>] [JUSTID]

    В команде:

    • <минимальное время> — неотрицательный порог с последней доставки в миллисекундах.
    • COUNT — задать положительное максимальное число передаваемых записей. По умолчанию — 100.
    • JUSTID — ограничить данные переданных записей их идентификаторами.

    Вернутся позиция для продолжения обхода, переданные записи и идентификаторы удаленных из потока записей, ссылки на которые были очищены из PEL. У переданных записей время с последней доставки будет сброшено. Если JUSTID не указан, счетчик доставок будет увеличен.

  7. Повторяйте вызов с полученной позицией вместо 0-0, сохраняя ключ, группу, потребителя и порог времени:

    text
    XAUTOCLAIM <имя ключа> <имя группы> <имя потребителя> <минимальное время> <следующий идентификатор> [COUNT <количество>] [JUSTID]

    Текущий обход завершится при возвращенной позиции 0-0. Пустой список переданных записей при другой позиции не означает завершение. Нулевой курсор не означает, что PEL пуст. Неподходящие по времени записи могут стать доступными для передачи позже.

  8. После успешной обработки подтвердите переданные записи.

Чтобы изменить позицию выдачи новых для группы записей, используйте команду XGROUP SETID.

Внимание

Перемещение позиции назад может привести к повторной выдаче записей, а вперед — к пропуску еще не выданных. Команда не выполнит повторную обработку самостоятельно и не очистит список неподтвержденных записей.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы и ее текущую позицию last-delivered-id.

  6. Выберите новую позицию. При необходимости получите идентификатор записи.

  7. Выполните команду:

    text
    XGROUP SETID <имя ключа> <имя группы> <идентификатор>

    В команде <идентификатор> — новая последняя позиция выдачи. Значение 0-0 вернет позицию к началу, а $ переместит ее к последней записи потока.

    Команда изменит позицию и вернет OK. Последующие вызовы XREADGROUP с > будут получать записи с большими идентификаторами. Содержимое потока и существующий PEL сохранятся. Если поток или группа отсутствует, команда вернет ошибку.

Чтобы удалить неиспользуемого потребителя группы, используйте команду XGROUP DELCONSUMER.

Внимание

После удаления потребителя его неподтвержденные записи нельзя будет передать через обычное перераспределение PEL. Сначала завершите их обработку или передайте другому потребителю. Не подтверждайте необработанные записи ради удаления потребителя.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы и узнайте имя потребителя.

  6. Проверьте его показатель pending. Если он не равен 0, передайте неподтвержденные записи другому потребителю или подтвердите уже успешно обработанные.

  7. Убедитесь, что приложение больше не использует имя удаляемого потребителя и его pending равен 0.

  8. Выполните команду:

    text
    XGROUP DELCONSUMER <имя ключа> <имя группы> <имя потребителя>

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

    Записи потока сохранятся. При последующей выдаче данных через XREADGROUP с тем же именем потребитель сможет появиться снова.

Чтобы удалить группу потребителей, используйте команду XGROUP DESTROY.

Внимание

Группа будет удалена даже при наличии активных потребителей и неподтвержденных записей. Состояние обработки этой группы будет потеряно.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте имя группы.

  6. Проверьте неподтвержденные записи и потребителей. Выполняйте удаление, когда приложения больше не используют эту группу.

  7. Выполните команду:

    text
    XGROUP DESTROY <имя ключа> <имя группы>

    Команда удалит группу вместе с ее потребителями и PEL и вернет 1. Если группы нет, вернется 0. Ключ потока, его записи и другие группы сохранятся.

    Повторное создание группы с тем же именем не восстановит прежний PEL и состояние потребителей.

Чтобы удалить старые записи по количеству или идентификатору, используйте команду XTRIM.

Внимание

Обрезка не будет ждать обработки записей всеми группами. Она может удалить еще не выданные или не подтвержденные записи, а ссылки на удаленные записи останутся в PEL. Пустой PEL не подтверждает, что все записи уже обработаны. Часть записей могла еще не выдаваться группе.

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Просмотрите записи и выберите границу хранения. Если используются группы, проверьте их состояние и неподтвержденные записи.

  6. Выполните нужную команду:

    • По количеству
    • По идентификатору
    text
    XTRIM <имя ключа> MAXLEN [= | ~] <количество>

    В команде <количество> — неотрицательное число записей, которые нужно оставить. При точной обрезке будут сохранены не более этого количества самых новых записей.

    Без оператора или с = обрезка будет точной. С ~ будут удаляться целые внутренние блоки, поэтому часть лишних записей может остаться. Только для режима ~ можно добавить LIMIT <лимит> — ограничение объема удаления за вызов. 0 снимет это ограничение.

    Команда вернет количество удаленных записей. Если удалять нечего или ключ отсутствует, вернется 0. При удалении всех записей сам поток сохранится.

Чтобы удалить отдельные записи потока, используйте команду XDEL.

Внимание

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

  1. Подключитесь к кластеру.

  2. Переключитесь на нужную базу данных.

  3. Узнайте имя ключа с помощью итеративного обхода.

  4. Проверьте тип значения ключа. Для существующего ключа должно вернуться значение stream.

  5. Узнайте идентификаторы удаляемых записей.

  6. Выполните команду:

    text
    XDEL <имя ключа> <идентификатор 1> [<идентификатор 2> ...]

    Команда удалит указанные записи и вернет количество удаленных записей. Несуществующие идентификаторы будут пропущены. Если ключ не существует, вернется 0. При удалении последней записи ключ потока сохранится.

    Удаление не заменит подтверждение обработки. Если запись еще учтена в PEL, ее идентификатор может остаться там без содержимого.