Продолжаем серию постов про Кафку. Сегодня расскажу про консьюмеров и связанную с ними историю про крит на проде 😱
📝 Консьюмеры — это те, кто читает сообщения из Кафки. Консьюмеры обычно состоят в консьюмер-группе (consumer group).
Как мы ранее разбирали, топики в Кафке разделены на партиции. Рассмотрим топик
orders с 4-мя партициями, из которого нам нужно читать. ➡️ Мы создали консьюмер группу
group-A, в ней один консьюмер. Он будет сам читать все 4 партиции. ➡️ Если мы добавим второй консьюмер, каждому достанется по 2 партиции.
➡️ Если мы сделаем четыре консьюмера в группе, то каждому достанется по одной партиции.
➡️ Если мы добавим пятого консьюмера, он будет простаивать, потому что партиций на него не хватило.
❗️Два консьюмера из одной консьюмер-группы не могут одновременно читать одну и ту же партицию.
Поэтому число партиций определяет максимальный параллелизм, с которым мы можем обрабатывать сообщения.
🤔 Но что если кому-то еще нужно читать данные из топика
orders?Всё в порядке, просто заводим другую консьюмер-группу
group-B со своими консьюмерами и спокойно читаем те же самые сообщения из orders абсолютно независимо от первой консьюмер-группы.📝 Если кто-то из консьюмеров выходит из строя или наоборот в группу добавляется новый консьюмер, то происходит переназначение партиций консьюмерам — ребалансировка.
Представим, что мы читаем из
orders, обрабатываем заказы. Тут один из консьюмеров отваливается. В этом случае произойдет ребалансировка, и все партиции будут поделены между оставшимися консьюмерами. Для распределения есть разные стратегии.Ребалансировку обычно видно в логах приложения с консьюмером (если все настроено стандартно).
🧐 Но как после ребалансировки свеженазначенный на партицию консьюмер знает, откуда ему продолжать чтение? Ведь сообщения в Кафке не удаляются после прочтения — они хранятся до истечения срока хранения (retention), независимо от того, вычитаны они или нет.
Для этого все консьюмеры регулярно отправляют брокеру информацию о том, какое последнее сообщение они прочитали. Это называется коммитом оффсета (offset commit) или по-русски фиксацией смещения.
Информация о смещениях для каждой партиции сохраняется в специальном служебном топике
__consumer_offsets. Поэтому свеженазначенный консьюмер будет знать, откуда продолжать чтение из партиции: с последнего закомиченного оффсета.🌟Способы коммитить оффсет
🟡 автоматически (
enable.auto.commit = true) — консьюмер без вмешательства разработчика будет коммитить последнее прочитанное смещение каждые 5 сек (по-умолчанию, можно изменить). В этом случае возможна ситуация, что одни и те же сообщения будут прочитаны и обработаны повторно.Например:
00:00 консьюмер зафиксировал оффсет 100
00:03 консьюмер прочитал и обработал несколько сообщений 101, 102, 103
00:04 консьюмер вышел из строя, не успев закоммитить оффсет
00:10 подключается новый консьюмер, получает информацию, что последний закомиченный оффсет был 100 и повторно вычитывает сообщения 101, 102, 103
❗️Поэтому очень важно, чтобы консьюмеры были идемпотентными. То есть умели определять дубликаты сообщений.
🟡 коммитить оффсет явно в коде после каждого прочитанного сообщения
🤔 Можно ли начать читать не с последнего закомиченного оффсета, а прочитать более ранние сообщения?
Можно.
➡️Первый вариант: делаем новую консьюмер-группу, и там начинаем жизнь с чистого листа: можно прочитать всё с самого старого из хранящихся сообщений или наоборот начать с самого свежего сообщения.
➡️ Второй вариант: можно самим явно задать смещение, с которого начать читать.
📝 Кафка выбирает одного из брокеров как group coordinator для каждой консьюмер-группы. Координатор группы:
🟢 отслеживает, кто входит в группу,
🔵запускает ребалансировку при необходимости,
🟣принимает и сохраняет коммиты оффсетов,
🟡следит за тем, жив ли каждый консьюмер
😼 Как координатор определяет, что консьюмер жив, а не отвалился?
✅ консьюмеры должны в фоне слать брокеру контрольные сигналы (heartbeats)
✅ консьюмеры должны периодически делать запрос новых сообщений (poll)
#kafka
Продолжение ⬇️
