Перейти к содержанию

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-реплик.

Пример:

Topic orders
  partition-0
    leader   -> broker-1
    follower -> broker-2
    follower -> broker-3

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 -> leader partition -> followers replicate

Это важно, потому что:

  • 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.

Пример:

replication.factor = 3
ISR = {broker-1, broker-2, broker-3}

Если broker-3 начал сильно лагать:

ISR = {broker-1, broker-2}

Тогда подтверждение acks=all уже будет идти по текущему ISR, а не по полному исходному набору реплик.


Commit и High Watermark

Запись в Kafka важно различать на двух уровнях:

  1. producer получил или не получил ack;
  2. запись committed и видна consumer'ам.

Упрощенно:

  • leader append'ит запись в свой лог;
  • followers из ISR догоняют leader;
  • когда запись считается committed, consumer может ее безопасно увидеть.

Практически это означает:

  • consumer не должен видеть "сырые" хвостовые записи, которые могут потеряться при failover;
  • consumer читает committed часть лога.

Эта граница committed-части логически соответствует high watermark.


Как запись проходит через кластер

Сценарий:

  1. Producer знает bootstrap.servers.
  2. Producer получает metadata о topic и лидерах partition.
  3. Producer выбирает partition.
  4. Producer отправляет record лидеру partition.
  5. Leader пишет запись в локальный append-only log.
  6. Followers подтягивают запись у leader.
  7. Когда условия commit выполнены, запись становится доступна consumer'ам.
  8. Producer получает ack в зависимости от acks.

Схема:

Producer
  -> broker leader
  -> local append
  -> follower replication
  -> commit
  -> ack / visibility to consumers

Как чтение проходит через кластер

Стандартный consumer:

  1. входит в consumer group;
  2. получает назначенные partition;
  3. отправляет fetch request лидерам этих partition;
  4. читает данные, начиная с нужного offset;
  5. двигает свою позицию;
  6. коммитит 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 = 3
  • min.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 вместе определяют реальные гарантии.