Группы потребителей и подписки на топики#
Чтобы позволить пулу процессов разделить работу по потреблению и обработке записей, Corax использует концепцию групп потребителей. При этом процессы могут как выполняться на одной машине, так и быть распределены по нескольким, чтобы обеспечить масштабируемость и отказоустойчивость.
Все экземпляры потребителей, использующие один и тот же group.id, будут входить в одну и ту же группу потребителей.
Каждый потребитель в группе может динамически задавать список топиков, на которые он хочет подписаться, с помощью API подписки subscribe. Corax доставит каждое сообщение в топиках, на которые подписались потребители, одному процессу в каждой группе потребителей. Это достигается за счет балансировки разделов между всеми участниками группы потребителей таким образом, чтобы каждый раздел был назначен ровно одному потребителю в группе. Таким образом, если существует топик с четырьмя разделами и группа потребителей с двумя процессами, каждый процесс будет потреблять из двух партиций.
Членство в группе потребителей поддерживается динамически: в случае сбоя процесса назначенные ему разделы будут переназначены другим потребителям в той же группе. Аналогично, если новый потребитель присоединится к группе, разделы будут перенесены с существующих потребителей на новых. Это называется перебалансировкой группы и более подробно описывается ниже.
Перебалансировка групп также применяется при добавлении новых разделов в один из топиков, у которого есть подписчики, или при создании нового топика, соответствующего регулярному выражению подписки. Группа автоматически обнаружит новые разделы с помощью периодического обновления метаданных и назначит их участникам группы.
Группу потребителей можно представлять как одного логического подписчика, который состоит из нескольких процессов. Как система с несколькими подписчиками, Corax поддерживает наличие любого количества групп потребителей для любого топика без дублирования данных (дополнительные потребители на самом деле не требуют большого количества ресурсов).
Это обобщение функциональности, распространенной в системах обмена сообщениями. Чтобы получить семантику, аналогичную очереди в традиционной системе обмена сообщениями, все процессы могут входить в одну группу потребителей, и, следовательно, доставка записей будет сбалансирована по группе, как в очереди. Однако в отличие от традиционной системы обмена сообщениями, у вас может быть несколько таких групп. Чтобы получить семантику, аналогичную pub-sub в традиционной системе обмена сообщениями, можно выделить каждому процессу свою собственную группу потребителей, и тогда каждый процесс будет подписываться на все записи, опубликованные в топике.
Кроме того, когда происходит автоматическое переназначение группы, потребители могут получать уведомления через слушателя ConsumerRebalanceListener, что позволяет им завершить необходимые операции на уровне приложения (например, очистку состояния, фиксацию смещения вручную и так далее).
Потребитель также может вручную назначить определенные разделы (аналогично более старому «простому» потребителю) с помощью метода assign(Collection). В этом случае динамическое назначение разделов и координация групп потребителей будут отключены.