XREADGROUP
XREADGROUP
XREADGROUP GROUP group consumer [COUNT count] [BLOCK milliseconds] [NOACK] STREAMS key [key ...] id [id ...]
- Доступно с версии:
- 5.0.0
- Сложность по времени:
- Для каждого упомянутого потока: O(M), где M — количество возвращённых элементов. Если M является константой (например, всегда запрашиваются первые 10 элементов с COUNT), можно считать его O(1). С другой стороны, при блокировке XREADGROUP, XADD будет тратить время O(N) для обслуживания N клиентов, заблокированных на потоке, ожидая новых данных.
- Категории ACL:
-
@write,@stream,@slow,@blocking,
Команда XREADGROUP — это специальная версия команды XREAD с поддержкой групп потребителей. Вероятно, вам нужно понять команду XREAD перед прочтением этой страницы, чтобы она имела смысл.
Кроме того, если вы новичок в потоках, мы рекомендуем прочитать наше введение в Redis Streams. Убедитесь, что вы поняли концепцию группы потребителей в этом введении, чтобы понять, как работает эта команда.
Группы потребителей за 30 секунд
Разница между этой командой и обычной командой XREAD заключается в том, что эта команда поддерживает группы потребителей.
Без групп потребителей, используя только XREAD, все клиенты получают все записи, поступающие в поток. Вместо этого, используя группы потребителей с XREADGROUP, можно создать группы клиентов, которые потребляют разные части сообщений, поступающих в данный поток. Например, если поток получает новые записи A, B и C, и есть два потребителя, читающие через группу потребителей, один клиент получит, например, сообщения A и C, а другой — сообщение B и так далее.
В рамках группы потребителей каждый потребитель (то есть клиент, потребляющий сообщения из потока) должен идентифицироваться уникальным именем потребителя. Это просто строка.
Одним из гарантий групп потребителей является то, что каждый потребитель может видеть только историю сообщений, которые были ему доставлены, поэтому сообщение принадлежит только одному владельцу. Однако существует специальная функция, называемая «заявлением на сообщение», которая позволяет другим потребителям заявлять о сообщениях в случае невосстановимого сбоя потребителя. Для реализации таких семантик, группы потребителей требуют явного подтверждения сообщений, успешно обработанных потребителем, через команду XACK. Это необходимо, потому что поток будет отслеживать для каждой группы потребителей, кто обрабатывает какое сообщение.
Вот как понять, нужно ли использовать группу потребителей или нет:
- Если у вас есть поток и несколько клиентов, и вы хотите, чтобы все клиенты получали все сообщения, вам не нужна группа потребителей.
- Если у вас есть поток и несколько клиентов, и вы хотите, чтобы поток был разделен или фрагментирован между вашими клиентами, так что каждый клиент получит подмножество сообщений, поступающих в поток, вам нужна группа потребителей.
Отличия между XREAD и XREADGROUP
С точки зрения синтаксиса команды почти одинаковы, но XREADGROUP требует специального и обязательного параметра:
GROUP <group-name> <consumer-name>
Имя группы — это просто имя группы потребителей, связанной с потоком. Группа создается с помощью команды XGROUP. Имя потребителя — это строка, используемая клиентом для идентификации себя в группе. Потребитель автоматически создается в группе потребителей при первом обнаружении. Разные клиенты должны выбирать разные имена потребителей.
При чтении с помощью XREADGROUP, сервер будет запоминать, что данное сообщение было вам доставлено: сообщение будет храниться в группе потребителей в списке ожидаемых записей (PEL), представляющем собой список идентификаторов сообщений, которые были доставлены, но ещё не подтверждены.
Клиент должен подтвердить обработку сообщения с помощью XACK, чтобы ожидаемая запись была удалена из PEL. PEL можно просмотреть с помощью команды XPENDING.
Подкоманда NOACK может использоваться для предотвращения добавления сообщения в PEL в тех случаях, когда надёжность не является требованием, и случайная потеря сообщения приемлема. Это эквивалентно подтверждению сообщения при его чтении.
Идентификатор для указания в параметре STREAMS при использовании XREADGROUP может быть одним из следующих двух:
- Специальный идентификатор
>, который означает, что потребитель хочет получить только сообщения, которые никогда не были доставлены ни одному другому потребителю. Это просто означает, дайте мне новые сообщения. - Любой другой идентификатор, то есть 0 или любой другой допустимый идентификатор или неполный идентификатор (только часть миллисекундного времени), будет иметь эффект возвращения записей, ожидающих потребителя, отправляющего команду с идентификаторами, большими, чем предоставленный. Иными словами, если идентификатор не
>, то команда просто позволит клиенту получить свои ожидающие записи: сообщения, которые были доставлены ему, но ещё не подтверждены. Обратите внимание, что в этом случае обаBLOCKиNOACKигнорируются.
Как и XREAD, команда XREADGROUP может использоваться в блокирующем режиме. В этом отношении нет различий.
Что происходит, когда сообщение доставляется потребителю?
Два события:
- Если сообщение никогда не доставлялось никому, то есть, если речь идёт о новом сообщении, создаётся PEL (Pending Entries List).
- Если вместо этого сообщение уже было доставлено этому потребителю, и он просто снова запрашивает то же сообщение, то счётчик последней доставки обновляется до текущего времени, а количество доставок увеличивается на единицу. Вы можете получить доступ к этим свойствам сообщения с помощью команды
XPENDING.
Пример использования
Обычно команда используется для получения новых сообщений и их обработки. В псевдокоде:
WHILE true
entries = XREADGROUP GROUP $GroupName $ConsumerName BLOCK 2000 COUNT 10 STREAMS mystream >
if entries == nil
puts "Timeout... try again"
CONTINUE
end
FOREACH entries AS stream_entries
FOREACH stream_entries as message
process_message(message.id,message.fields)
# ACK the message as processed
XACK mystream $GroupName message.id
END
END
END
Таким образом, пример кода потребителя получит только новые сообщения, обработает их и подтвердит их с помощью XACK. Однако пример кода выше не является полным, так как он не обрабатывает восстановление после сбоя. Что произойдёт, если произойдёт сбой в середине обработки сообщений, так наши сообщения останутся в списке ожидающих записей, поэтому мы можем получить доступ к нашей истории, предоставив XREADGROUP первоначально идентификатор 0 и выполнив тот же цикл. После предоставления идентификатора 0 ответ — пустой набор сообщений, мы знаем, что мы обработали и подтвердили все ожидающие сообщения: мы можем начать использовать > в качестве идентификатора, чтобы получить новые сообщения и присоединиться к потребителям, обрабатывающим новые вещи.
Чтобы увидеть, как команда фактически отвечает, пожалуйста, проверьте страницу команды XREAD.
Что происходит, когда ожидающее сообщение удаляется?
Записи могут быть удалены из потока из-за обрезки или явных вызовов XDEL в любое время. По дизайну, Redis не предотвращает удаление записей, присутствующих в PEL потока. В этом случае PEL сохраняет идентификаторы удалённых записей, но фактический payload записи больше недоступен. Поэтому при чтении таких записей PEL Redis вернёт значение null вместо их соответствующих данных.
Пример:
> XADD mystream 1 myfield mydata
"1-0"
> XGROUP CREATE mystream mygroup 0
OK
> XREADGROUP GROUP mygroup myconsumer STREAMS mystream >
1) 1) "mystream"
2) 1) 1) "1-0"
2) 1) "myfield"
2) "mydata"
> XDEL mystream 1-0
(integer) 1
> XREADGROUP GROUP mygroup myconsumer STREAMS mystream 0
1) 1) "mystream"
2) 1) 1) "1-0"
2) (nil)
Возвращаемое значение
Массивный ответ, конкретно:
Команда возвращает массив результатов: каждый элемент возвращаемого массива — это массив, состоящий из двух элементов, содержащих имя ключа и записи, возвращённые для этого ключа. Возвращённые записи — полные записи потока, содержащие идентификаторы и список всех полей и значений. Гарантируется, что поля и значения возвращаются в том же порядке, в котором они были добавлены с помощью XADD.
При использовании BLOCK, при истечении времени ожидания возвращается нулевой ответ.
Настоятельно рекомендуется прочитать введение в Redis Streams, чтобы лучше понять поведение и семантику потоков в целом.
© 2006–2022 Salvatore Sanfilippo
Licensed under the Creative Commons Attribution-ShareAlike License 4.0.
https://redis.io/commands/xreadgroup/