XREAD
XREAD
XREAD [COUNT count] [BLOCK milliseconds] STREAMS key [key ...] id [id ...]
- Доступно с версии:
- 5.0.0
- Сложность по времени:
- Категории ACL:
-
@read,@stream,@slow,@blocking,
Считывает данные из одного или нескольких потоков, возвращая только записи с идентификатором, большим, чем последний полученный идентификатор, сообщённый вызывающей стороной. Эта команда имеет опцию блокировки, если элементы недоступны, аналогично BRPOP или BZPOPMIN и другим.
Обратите внимание, что перед чтением этой страницы, если вы новичок в потоках, рекомендуется прочитать введение в Redis Streams.
Неблокирующее использование
Если опция BLOCK не используется, команда является синхронной и может считаться чем-то похожей на XRANGE: она вернёт диапазон элементов внутри потоков, однако у неё есть два ключевых отличия по сравнению с XRANGE, даже если мы рассмотрим только синхронное использование:
- Эту команду можно вызвать с несколькими потоками, если мы хотим одновременно читать из нескольких ключей. Это ключевая особенность
XREAD, потому что особенно при блокировке с помощью BLOCK, возможность прослушивания по нескольким ключам с одним соединением — важная функция. - В то время как
XRANGEвозвращает элементы в диапазоне идентификаторов,XREADболее подходит для потребления потока, начиная с первой записи, которая больше, чем любая другая запись, которую мы видели до сих пор. Поэтому то, что мы передаём вXREAD, это для каждого потока идентификатор последнего элемента, который мы получили из этого потока.
Например, если у меня есть два потока mystream и writers, и я хочу читать данные из обоих потоков, начиная с первого элемента, который они содержат, я мог бы вызвать XREAD следующим образом.
Примечание: в примере используется опция COUNT, поэтому для каждого потока вызов вернёт максимум два элемента на поток.
> XREAD COUNT 2 STREAMS mystream writers 0-0 0-0
1) 1) "mystream"
2) 1) 1) 1526984818136-0
2) 1) "duration"
2) "1532"
3) "event-id"
4) "5"
5) "user-id"
6) "7782813"
2) 1) 1526999352406-0
2) 1) "duration"
2) "812"
3) "event-id"
4) "9"
5) "user-id"
6) "388234"
2) 1) "writers"
2) 1) 1) 1526985676425-0
2) 1) "name"
2) "Virginia"
3) "surname"
4) "Woolf"
2) 1) 1526985685298-0
2) 1) "name"
2) "Jane"
3) "surname"
4) "Austen"
Опция STREAMS является обязательной и ДОЛЖНА быть последней опцией, так как такая опция получает аргументы переменной длины в следующем формате:
STREAMS key_1 key_2 key_3 ... key_N ID_1 ID_2 ID_3 ... ID_N
Итак, мы начинаем со списка ключей, а затем продолжаем со всеми связанными идентификаторами, представляющими последний полученный идентификатор для этого потока, чтобы вызов возвращал только идентификаторы больше, чем у того же потока.
Например, в примере выше последний элемент, который мы получили для потока mystream имеет идентификатор 1526999352406-0, а для потока writers — идентификатор 1526985685298-0.
Чтобы продолжить итерацию по двум потокам, я вызову:
> XREAD COUNT 2 STREAMS mystream writers 1526999352406-0 1526985685298-0
1) 1) "mystream"
2) 1) 1) 1526999626221-0
2) 1) "duration"
2) "911"
3) "event-id"
4) "7"
5) "user-id"
6) "9488232"
2) 1) "writers"
2) 1) 1) 1526985691746-0
2) 1) "name"
2) "Toni"
3) "surname"
4) "Morrison"
2) 1) 1526985712947-0
2) 1) "name"
2) "Agatha"
3) "surname"
4) "Christie"
И так далее. В конечном итоге вызов не вернёт ни одного элемента, а вернёт только пустой массив, тогда мы знаем, что больше ничего нет для извлечения из потока (и нам придётся повторить операцию, поэтому эта команда также поддерживает режим блокировки).
Неполные идентификаторы
Использование неполных идентификаторов допустимо, как и в XRANGE. Однако в данном случае последовательная часть идентификатора, если отсутствует, всегда интерпретируется как ноль, поэтому команда:
> XREAD COUNT 2 STREAMS mystream writers 0 0
точно эквивалентна
> XREAD COUNT 2 STREAMS mystream writers 0-0 0-0
Блокировка для данных
В своей синхронной форме команда может получать новые данные, пока доступны дополнительные элементы. Однако в какой-то момент нам придётся подождать, пока производители данных не воспользуются XADD для добавления новых записей в потоки, которые мы потребляем. Чтобы избежать опроса с фиксированным или адаптивным интервалом, команда может заблокироваться, если не смогла вернуть какие-либо данные, согласно указанным потокам и идентификаторам, и автоматически разблокируется, как только один из запрошенных ключей примет данные.
Важно понимать, что эта команда распределяет данные всем клиентам, которые ожидают того же диапазона идентификаторов, поэтому каждый потребитель получит копию данных, в отличие от того, что происходит при использовании операций блокировки извлечения из списков.
Для блокировки используется опция BLOCK вместе с количеством миллисекунд, на которое мы хотим заблокироваться, прежде чем истечёт таймаут. Обычно команды блокировки Redis принимают таймауты в секундах, однако эта команда принимает таймаут в миллисекундах, даже если обычно у сервера есть разрешение таймаута около 0,1 секунды. В некоторых случаях это позволяет заблокироваться на более короткое время, и при улучшении внутренней реализации сервера, возможно, разрешение таймаутов улучшится.
Когда опция BLOCK передаётся, но есть данные для возврата по крайней мере в одном из переданных потоков, команда выполняется синхронно, точно так же, как если бы опция BLOCK отсутствовала.
Вот пример вызова с блокировкой, где команда впоследствии возвращает null-ответ, потому что таймаут истек, и новые данные не поступили:
> XREAD BLOCK 1000 STREAMS mystream 1526999626221-0 (nil)
Специальный идентификатор $
Иногда при блокировке мы хотим получить только записи, добавленные в поток с помощью XADD с момента блокировки. В таком случае нас не интересует история уже добавленных записей. Для этого случая нам нужно проверить идентификатор верхнего элемента потока и использовать этот идентификатор в команде XREAD. Это неэффективно и требует вызова других команд, поэтому вместо этого можно использовать специальный идентификатор $ для сигнализации потоку о том, что мы хотим только новые данные.
Очень важно понять, что вы должны использовать идентификатор $ только для первого вызова XREAD. Позже идентификатор должен быть идентификатором последнего отчётного элемента в потоке, иначе вы можете пропустить все записи, добавленные между ними.
Вот как выглядит типичный вызов XREAD в первой итерации потребителя, который хочет потреблять только новые записи:
> XREAD BLOCK 5000 COUNT 100 STREAMS mystream $
После получения ответов следующий вызов будет выглядеть примерно так:
> XREAD BLOCK 5000 COUNT 100 STREAMS mystream 1526999644174-3
И так далее.
Как обслуживаются несколько клиентов, заблокированных на одном потоке
Операции блокировки извлечения из списков или наборов с упорядоченными множествами имеют поведение pop. В принципе, элемент удаляется из списка или упорядоченного множества, чтобы он был возвращён клиенту. В этом случае вы хотите, чтобы элементы потреблялись справедливо, в зависимости от момента, когда клиенты заблокировались на данном ключе. Обычно Redis использует семантику FIFO в таких случаях.
Однако обратите внимание, что с потоками это не проблема: записи потока не удаляются из потока при обслуживании клиентов, поэтому каждый ожидающий клиент будет обслуживаться, как только команда XADD предоставит данные в поток.
Возврат
Ответ в виде массива, конкретно:
Команда возвращает массив результатов: каждый элемент возвращаемого массива представляет собой массив, состоящий из двух элементов, содержащих имя ключа и отчётные записи для этого ключа. Отчётные записи — это полные записи потока, имеющие идентификаторы и список всех полей и значений. Поля и значения гарантированно возвращаются в том же порядке, в котором они были добавлены командой XADD.
При использовании BLOCK, при истечении таймаута возвращается null-ответ.
Для лучшего понимания общего поведения и семантики потоков настоятельно рекомендуется ознакомиться с введением в Redis Streams.
© 2006–2022 Salvatore Sanfilippo
Licensed under the Creative Commons Attribution-ShareAlike License 4.0.
https://redis.io/commands/xread/