Продолжаем рассматривать особенности share groups, которые обеспечивают работу с сообщениями как в очередях.
➡️ Чтение каждого сообщения подтверждается индивидуально
В случае consumer groups чтение фиксируется коммитом оффсета для каждой партиции. Если для этой консьюмер-группы закоммичен оффсет на сообщении 90, значит все сообщения до 90 прочитаны.
В share группе появляется понятие окна чтения (in-flight records) и статусная модель сообщения.
Статусы:
🌟 Available — сообщение доступно для чтения консьюмером
🌟 Acquired — сообщения было взято в обработку (залочено) конкретным консьюмером (время, на которое сообщение залочено, настраивается параметром
share.record.lock.duration.ms, по-умолчанию 30 сек)🌟 Acknowledged — сообщение обработано консьюмером (можно сказать, что это аналог коммита оффсета, но лишь для одного сообщения)
🌟 Archived — сообщение признано архивным, больше недоступно для чтения никаким консьюмером в этой share группе.
➡️ Лимит количества попыток доставки сообщения
Для share группы настраивается параметр
group.share.delivery.attempt.limit (по-умолчанию 5) — это лимит попыток доставки. Для каждого сообщения Kafka ведет счетчик. Лимит позволяет бороться с невалидными сообщениями, чтобы они не доставлялись бесконечно.Жизненный цикл сообщение в share группе:
1️⃣ Сначала сообщение находится в статусе Available. Его забирает на обработку какой-то из консьюмеров. Счетчик попыток доставки увеличивается на 1. Если у сообщения еще не исчерпан лимит попыток обработки, то его статус становится Acquired. Если лимит попыток исчерпан, то оно становится Archived.
2️⃣ Если консьюмер обработал сообщение, то оно переводится в статус Acknowledged. Также консьюмер может зареджектить сообщение (например, если оно невалидное). Тогда оно становится Archived. Если консьюмер не успел обработать сообщение за время лока (30 сек), то оно снова становится Available.
3️⃣ Со временем, когда окно чтения сдвигается, все Acknowledged из прошлого окна становятся Archived.
Важно, что эти статусы работают в рамках конкретной share group. В одной share group это сообщение уже Archived, а в другой - Available.
Топики и партиции остаются теми же, что были. Один сервис может читать топик через consumer группу, а другой - читать этот же топик через share группу. Продюсер вообще не знает, кто там и как читает. Его дело в партиции записать.
Share groups — это просто другой тип группы потребителей, использующий другую модель распределения записей и другой протокол подтверждения:
🟢 consumer groups читают топик как лог
🟢 share groups читают тот же топик как очередь
Когда может пригодиться новый подход с share groups?
Если нам нужна очередь, а не лог: сообщения - это независимые события, которые можно обрабатывать параллельно, порядок некритичен, при этом мы хотим динамически масштабировать их обработку. Хотим ускорить обработку - добавляем больше воркеров (консьюмеров).
Подробно почитать про нововведение можно в доке: KIP-932: Queues for Kafka
Очень круто, что Kafka так развивается! Конечно, фича новая, сейчас правят баги, но для будущих проектов считаю интересно иметь в виду.
#kafka