Партиции, офсеты и группы потребителей в Kafka
Kafka - де-факто лог индустрии, и её внутренности объясняют почти всё поведение: порядок, параллелизм, надёжность. Собес проверяет, понимаешь ли ты, что порядок только внутри партиции и почему acks=all - стандарт надёжности.
Стержень: топик это партиции; порядок гарантирован лишь внутри партиции, а ключ сообщения определяет, в какую партицию оно ляжет.
// Формулировки: «как Kafka гарантирует порядок?», «зачем ключ сообщения?», «что такое ISR и acks=all?».
Партиции и порядок
Топик состоит из партиций, и партиция это упорядоченный лог на диске. Порядок гарантирован только внутри партиции, глобального порядка по топику нет. Ключ сообщения хешируется в номер партиции, поэтому все события одной сущности (один заказ) ложатся в одну партицию и приходят по порядку.
Число партиций это потолок параллелизма консюмер-группы: одну партицию читает один консюмер группы. Менять число партиций потом больно - оно ломает распределение ключей.
// Отсюда ловушка перекоса: ключ с клиентом-гигантом создаёт горячую партицию, у которой растёт лаг, пока остальные простаивают.
- партиция
- упорядоченный лог; порядок только внутри неё
- ключ → партиция
- hash ключа фиксирует партицию и порядок сущности
Группы, репликация, надёжность
Consumer group делит партиции между членами; смерть или появление члена вызывает rebalance - паузу чтения. Статическое членство и cooperative-протокол смягчают ребалансы. Консюмеров в группе больше числа партиций держать бессмысленно - лишние простаивают.
Репликация: у партиции есть лидер (через него всё чтение-запись) и фолловеры; ISR (in-sync replicas) - реплики в синхроне. Стандарт надёжности записи - acks=all вместе с min.insync.replicas=2: запись подтверждается только когда её приняли синхронные реплики.
// acks=1 для критичных данных опасен: лидер подтвердил, но упал до репликации - данные потеряны при смене лидера.
- partition leader / ISR
- точка записи / реплики в синхроне
- acks=all
- подтверждение записи всеми синхронными репликами
Продюсер, диск, retention
Продюсер настраивают трейдофом: acks (0/1/all), батчинг (linger.ms + batch.size) меняет throughput против латентности, сжатие батча (zstd/lz4). Ключ даёт контроль партиционирования и порядка.
Kafka быстра, потому что работает последовательным диском: append в лог, zero-copy чтение, page cache - случайного доступа нет, отсюда и дешевизна хранения потока.
// Retention идёт по времени или размеру, независимо от прочтения. Compacted topic хранит последнее значение по ключу это changelog-семантика для состояний и CDC (change data capture).
- rebalance
- передел партиций между консюмерами группы
- log compaction
- хранение последнего значения на ключ (changelog)
Как отвечать: «Как Kafka гарантирует порядок и надёжность записи?»
Про порядок: Kafka гарантирует его только внутри партиции, глобального порядка по топику нет. Партиция это упорядоченный лог на диске, и чтобы события одной сущности шли по порядку, я задаю ключ сообщения - например, id заказа: ключ хешируется в номер партиции, и все события этого заказа попадают в одну партицию, где порядок сохранён. Если порядок не важен или ключ распределён равномерно - параллелизм растёт с числом партиций. Про надёжность записи: у партиции есть лидер и реплики-фолловеры, синхронные из них образуют ISR. Стандарт - ставить acks=all вместе с min.insync.replicas=2, тогда продюсер получает подтверждение только после того, как запись приняли синхронные реплики, и потеря лидера не теряет данные. acks=1 для критичного опасен - лидер мог подтвердить и упасть до репликации. То есть порядок держит ключ и партиция, надёжность - acks=all плюс ISR.
Почему это сильный ответ: порядок объяснён через партиции и ключ (не глобальный), надёжность - через ISR и acks=all с контрпримером acks=1, и связка с параллелизмом - механика, а не термины.
На чём валят
- −Ждать глобального порядка по топику - он только внутри партиции.
- −Ключ с перекосом (клиент-гигант) - горячая партиция, лаг при простое остальных.
- −acks=1 для критичных данных - потеря при смене лидера до репликации.
- −Консюмеров в группе больше, чем партиций - лишние стоят без работы.
- −Тяжёлая обработка в цикле poll без пауз/async - таймаут сессии, вечный ребаланс-шторм.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 10, остальные разбираются в тренажёре.
- Что происходит при ребалансировке (rebalance) consumer group?A)Перераспределение партиций по консьюмерам, короткая пауза чтенияB)Kafka физически перемещает сами данные партиций между брокерами кластера, освобождая местоC)Ребаланс происходит на каждое отдельное сообщение, обеспечивая идеально ровную нагрузкуD)Во время ребаланса потребление ускоряется, так как подключается больше консьюмеров разом
показать ответ и разбор
+A)Перераспределение партиций по консьюмерам, короткая пауза чтения// разбор: Когда консьюмер входит в группу, выходит или признан упавшим (не шлёт heartbeat), координатор запускает ребаланс: партиции переназначаются между оставшимися консьюмерами так, чтобы каждую читал ровно один. На время stop-the-world ребаланса потребление в группе приостанавливается — частые ребалансы (долгая обработка, короткий таймаут) бьют по throughput. Смягчают cooperative rebalancing и настройкой таймаутов/heartbeat.
- Как выбор ключа сообщения (message key) влияет на партиционирование в Kafka?A)Ключ сообщения не влияет на выбор партиции — Kafka раскидывает записи случайноB)Хеш ключа фиксирует партицию (порядок по ключу), но скос ключа → hot partitionC)Сообщения с одинаковым ключом Kafka намеренно раскидывает по разным партициям для балансаD)Ключ определяет, какой именно консьюмер прочитает сообщение, минуя механизм партиций
показать ответ и разбор
+B)Хеш ключа фиксирует партицию (порядок по ключу), но скос ключа → hot partition// разбор: Продюсер по умолчанию выбирает партицию как hash(key) % partitions, поэтому все сообщения с одним ключом попадают в одну партицию — это даёт порядок в пределах ключа (все события юзера по порядку). Обратная сторона: если один ключ доминирует (популярный пользователь), его партиция перегружается — hot partition, ограничивающая throughput. Ключ подбирают достаточно равномерным; сообщения без ключа раскидываются по партициям равномерно (но без гарантии порядка).
- Что задаёт политика retention топика Kafka?A)Retention задаёт, скольким консьюмерам разрешено одновременно читать топик до его блокировкиB)Как долго/сколько данных хранить в логе (по времени или размеру), прежде чем удалять старые сегментыC)Retention определяет максимальный размер одного отдельного сообщения, допустимого в топикеD)После истечения retention Kafka физически удаляет весь топик вместе с его конфигурацией целиком
показать ответ и разбор
+B)Как долго/сколько данных хранить в логе (по времени или размеру), прежде чем удалять старые сегменты// разбор: Retention определяет, сколько лог живёт: по времени (retention.ms, например 7 дней) или по размеру (retention.bytes). Старые сегменты за границей удаляются. Это важно для replay: окно переигровки не больше retention — нельзя перечитать то, что уже удалено. Больше retention — дороже диск, но шире окно восстановления/бэкфилла. Компактированные топики живут по другой политике (последнее значение на ключ). Retention подбирают под сценарии переигровки и стоимость хранения.
- Что такое брокер (broker) в Kafka?A)клиентская библиотека, отправляющая сообщения в топикB)сервер кластера, хранящий партиции и отдающий их клиентамC)процесс, преобразующий сообщения между форматамиD)координатор, распределяющий партиции между консьюмерами
показать ответ и разбор
+B)сервер кластера, хранящий партиции и отдающий их клиентам// разбор: Брокер — узел кластера: он держит на диске партиции топиков, принимает записи от продюсеров и отдаёт сообщения консьюмерам. Кластер состоит из нескольких брокеров, партиции распределены между ними, и это даёт масштабирование и отказоустойчивость. Координация консьюмеров и выбор лидеров партиций — тоже работа брокеров, одному из которых достаётся роль контроллера.
- Зачем у партиции топика заводят несколько реплик?A)чтобы несколько консьюмеров читали её параллельноB)чтобы ускорить запись: продюсер пишет во все реплики сразуC)чтобы данные пережили отказ брокера, на котором лежит лидерD)чтобы хранить разные версии сообщений одного ключа
показать ответ и разбор
+C)чтобы данные пережили отказ брокера, на котором лежит лидер// разбор: Одна реплика партиции — лидер, остальные тянут с неё копию. Упал брокер с лидером — одна из синхронных реплик становится новым лидером, и данные не теряются. Параллельность чтения даёт не репликация, а число партиций: внутри группы одну партицию читает ровно один консьюмер. Продюсер тоже пишет только лидеру, а не всем репликам сразу.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.