TGViewer
Женя Янченко Женя Янченко @jane_yanchenko · 5.51K subscribers
Post #305 4.28K
Очереди в Кафке (продолжение)

Продолжаем рассматривать особенности 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
  • 👍 23
  • ❤ 18
  • 🔥 14
  • ❤‍🔥 3
  • 🐳 1
More from @jane_yanchenko
  1. Sep 21, 2026🔗 Подборка постов про Кафку Как обещала на стриме, собрала посты про Кафку в удобное огла…
  2. Sep 21, 2026🎞 Готова запись стрима про Кафку: https://youtu.be/2aRKsD-MWDA Большое спасибо всем, кто…
  3. Sep 16, 2026Сегодня стрим по Кафке в 19:00 Планируем не в формате доклада, а в формате вопрос-ответ, ч…
  4. Sep 16, 2026Post #413
  5. Sep 16, 2026Post #412
  6. Sep 16, 2026Post #411
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →