Kafka Architecture¶
Как устроен Kafka-кластер¶
В кластере Kafka обычно есть:
- несколько
broker-ов; - metadata layer для управления кластером;
- topics, разбитые на partition;
- replicas этих partition на разных broker'ах.
Упрощенно:
Kafka Cluster
broker-1
broker-2
broker-3
Topic orders.created
partition-0 -> replicas on broker-1, broker-2, broker-3
partition-1 -> replicas on broker-2, broker-3, broker-1
partition-2 -> replicas on broker-3, broker-1, broker-2
Broker¶
Broker — это узел Kafka, который:
- хранит partition'ы;
- принимает produce/fetch запросы;
- реплицирует данные;
- участвует в лидерстве partition'ов;
- сообщает metadata клиентам.
Важно:
- broker хранит не "весь topic целиком", а только те partition/replica, которые ему назначены;
- обычно partition'ов намного больше, чем broker'ов;
- нагрузка распределяется именно через partition'ы.
Controller и metadata¶
Кластеру нужен компонент, который управляет metadata:
- какие broker'ы живы;
- где лидер каждой partition;
- какие replicas входят в ISR;
- кто должен стать новым лидером при сбое.
В современной Kafka обычно используется KRaft:
- есть controller quorum;
- controller управляет metadata кластера;
- broker'ы получают актуальное состояние от controller'ов.
Исторически Kafka часто работала через ZooKeeper, поэтому в старых материалах можно увидеть описание через ZK.
Практический вывод:
- для понимания архитектуры важнее помнить не ZooKeeper vs KRaft;
- важнее помнить, что кластеру нужен metadata/control plane, который знает состояние broker'ов, partition'ов и лидеров.
Topic, partition и replica¶
Kafka хранит данные не просто "по topic", а по partition.
У каждой partition есть:
- один
leader; - ноль или больше
follower-реплик.
Пример:
replication.factor = 3 означает:
- у partition будет 3 копии данных;
- одна из них лидер;
- остальные follower'ы.
Важно:
- unit of scaling = partition;
- unit of replication = partition;
- unit of ordering = partition.
Leader и follower¶
Leader partition — это реплика, которая принимает запись и обслуживает стандартные клиентские чтения.
Follower replica:
- копирует данные у leader;
- держит такую же последовательность записей;
- может стать новым лидером при failover, если достаточно актуальна.
То есть путь записи выглядит так:
Это важно, потому что:
- producer пишет не "в topic вообще", а в leader конкретной partition;
- отказоустойчивость обеспечивается не магией topic'а, а репликацией partition.
ISR (In-Sync Replicas)¶
ISR — это набор реплик partition, которые считаются достаточно синхронными с leader.
Если follower:
- отстал слишком сильно;
- временно недоступен;
- не успевает догонять leader;
то его могут убрать из ISR.
Зачем это важно:
- commit в Kafka завязан именно на ISR;
- leader election обычно выбирает нового лидера из ISR;
acks=allработает относительно текущего ISR.
Пример:
Если broker-3 начал сильно лагать:
Тогда подтверждение acks=all уже будет идти по текущему ISR, а не по полному исходному набору реплик.
Commit и High Watermark¶
Запись в Kafka важно различать на двух уровнях:
- producer получил или не получил ack;
- запись committed и видна consumer'ам.
Упрощенно:
- leader append'ит запись в свой лог;
- followers из ISR догоняют leader;
- когда запись считается committed, consumer может ее безопасно увидеть.
Практически это означает:
- consumer не должен видеть "сырые" хвостовые записи, которые могут потеряться при failover;
- consumer читает committed часть лога.
Эта граница committed-части логически соответствует high watermark.
Как запись проходит через кластер¶
Сценарий:
- Producer знает
bootstrap.servers. - Producer получает metadata о topic и лидерах partition.
- Producer выбирает partition.
- Producer отправляет record лидеру partition.
- Leader пишет запись в локальный append-only log.
- Followers подтягивают запись у leader.
- Когда условия commit выполнены, запись становится доступна consumer'ам.
- Producer получает ack в зависимости от
acks.
Схема:
Producer
-> broker leader
-> local append
-> follower replication
-> commit
-> ack / visibility to consumers
Как чтение проходит через кластер¶
Стандартный consumer:
- входит в consumer group;
- получает назначенные partition;
- отправляет fetch request лидерам этих partition;
- читает данные, начиная с нужного offset;
- двигает свою позицию;
- коммитит offset при необходимости.
Важно:
- consumer не спрашивает "дай следующее непрочитанное сообщение вообще";
- consumer читает конкретный лог конкретной partition, начиная с конкретной позиции.
Failover: что происходит при падении broker¶
Сценарий 1. Падает follower¶
Если упал follower:
- лидер продолжает обслуживать запись и чтение;
- ISR может уменьшиться;
- кластер продолжает работать;
- durability запас уменьшается.
Сценарий 2. Падает leader¶
Если упал leader:
- controller замечает отказ;
- выбирается новый leader из актуальных ISR;
- producers/consumers обновляют metadata;
- после короткой паузы работа продолжается.
Важно:
- чем здоровее ISR, тем безопаснее failover;
- если ISR маленький, риски выше;
- лидер выбирается не произвольно, а из достаточно актуальных реплик.
Сценарий 3. Живых ISR не осталось¶
Это плохой сценарий.
Тогда возникает выбор между:
- consistency: ждать, пока поднимется подходящая реплика;
- availability: разрешить стать лидером не полностью синхронной реплике и рискнуть потерей части данных.
В Kafka этот компромисс связан с unclean.leader.election.enable.
Идея:
- включаешь unclean election -> выше шанс доступности, выше риск потери подтвержденных данных;
- выключаешь -> кластер предпочитает консистентность, но может стать временно недоступен на запись/чтение partition.
Для продовых критичных событий обычно предпочитают consistency.
Replication factor и min.insync.replicas¶
Это две разные вещи.
replication.factor¶
Сколько всего копий partition хранится в кластере.
Например:
replication.factor = 3- есть 1 leader + 2 follower
min.insync.replicas¶
Минимальное число ISR-реплик, которое нужно для успешной записи при acks=all.
Типичный продовый вариант:
replication.factor = 3min.insync.replicas = 2- producer использует
acks=all
Это дает хороший practical balance между durability и availability.
Почему partition так важны¶
Partition одновременно определяют:
- масштабирование producer/consumer throughput;
- parallelism;
- порядок сообщений;
- распределение данных по broker'ам;
- распределение лидерства по кластеру.
Отсюда важный вывод:
Выбор числа partition — это не косметика, а архитектурное решение.
Слишком мало partition:
- мало параллелизма;
- узкое место по throughput;
- один hot key может забивать topic.
Слишком много partition:
- растет нагрузка на metadata и coordination;
- больше файлов, больше сетевой и operational overhead;
- тяжелее rebalances и recovery.
Что запомнить¶
- Kafka масштабируется через partition.
- Отказоустойчивость Kafka строится на replication.
- У partition один leader и несколько follower'ов.
- Producer пишет в leader partition.
- Consumer читает committed log по offset.
- ISR — это набор реплик, которым сейчас доверяют для commit и failover.
replication.factor,acksиmin.insync.replicasвместе определяют реальные гарантии.