Эксплуатация Kafka: KRaft, мониторинг, тюнинг, обновление

Kafka в проде: KRaft (без ZooKeeper), ключевые метрики (consumer lag, under-replicated, latency), тюнинг producer и consumer, partition reassignment и rack awareness, rolling upgrade и квоты — продакшн-чеклист

Kafka несложно поднять и мучительно эксплуатировать вслепую: кластер молча накапливает under-replicated партиции, consumer lag растёт, а первым сигналом становится инцидент. Эта статья — про операционный слой: чем управляет KRaft-кворум контроллеров и чем управлял ZooKeeper до него, что мониторить, как тюнить producer и consumer не наугад, как переливать партиции между брокерами и почему это не то же самое, что ребаланс consumer group, и как квоты защищают кластер от одного шумного клиента.

Это шестая статья серии «Kafka: глубокое погружение» — операционная. Опирается на модель ISR и кворума контроллера из статьи про репликацию и на модель сегментов диска из статьи про хранение; явно разводит два похожих по звучанию, но разных по природе механизма — ребаланс consumer group (кто из потребителей читает партицию) и administrative reassignment партиций (на каком брокере физически лежит реплика), — которые путают чаще всего. Все числа ниже — из живого трёхброкерного KRaft-кластера. Все сценарии открыты в репозитории digital-cookbook — стенды kafka/go/ops (franz-go) и kafka/java/ops (kafka-clients), плюс bash-обвязка в kafka/ops/.

Эксплуатация Kafka: KRaft-кворум вместо ZooKeeper, метрики, reassignment и тюнинг

В статье

KRaft и эпоха ZooKeeper: что изменилось

До KRaft Kafka не была самодостаточной системой хранения собственных метаданных — она делегировала эту работу внешнему ансамблю ZooKeeper. ZooKeeper отвечал за три вещи одновременно: хранил метаданные кластера (список топиков, число партиций, конфиги, ACL) в виде дерева znode’ов; вёл выборы контроллера — единственного брокера, ответственного за административные решения (создание топиков, переизбрание лидеров партиций, обработку выхода брокеров из строя); и отслеживал liveness брокеров через ephemeral znode’ы — если брокер переставал продлевать свою сессию, ZooKeeper считал его мёртвым и уведомлял контроллер.

Схема работала, но с ростом кластеров упиралась в структурные ограничения, которые и стали формальным обоснованием KIP-500:

  • Две системы вместо одной. ZooKeeper — отдельный распределённый сервис со своим кворумом, своими бэкапами, своим мониторингом и своей версией развёртывания, независимой от Kafka. Каждый инцидент требовал понимания двух разных систем сразу.
  • Ограниченный масштаб метаданных. Znode’ы ZooKeeper не рассчитаны на десятки и сотни тысяч партиций — под сильной нагрузкой метаданных сама ZooKeeper-подсистема становилась узким местом кластера.
  • Медленный failover контроллера. При смене контроллера новый узел был обязан вычитать всё состояние метаданных из ZooKeeper заново, прежде чем начать работать, — на большом кластере это были секунды-десятки секунд простоя административных операций, а не мгновенное переключение.
  • Рассинхронизация двух источников правды. Метаданные в ZooKeeper и фактическое состояние брокеров могли временно разойтись просто потому, что это два разных процесса со своей задержкой репликации.

KRaft убирает внешнюю систему целиком: метаданные хранятся в самой Kafka как обычный реплицируемый лог — специальный внутренний топик __cluster_metadata — и распространяются через Raft-кворум контроллеров, встроенный прямо в брокеров. Роль controller может быть совмещена с ролью broker на одном и том же узле (combined, как в примере ниже) или вынесена на отдельные выделенные узлы (dedicated) — выбор влияет на то, что именно ломается при потере узла, к этому вернёмся в следующем разделе. Что это даёт на практике: одна система вместо двух, быстрый failover контроллера (новый лидер Raft-группы уже держит актуальный лог, вычитывать состояние заново из внешнего хранилища не нужно), кратно больший практический потолок числа партиций на кластер и единая непротиворечивая модель кворума вместо двух разных.

Таймлайн перехода занял не один релиз: KRaft получил статус production-ready (GA) в Kafka 3.3 (октябрь 2022, KIP-833); режим на основе ZooKeeper был официально помечен deprecated в Kafka 3.5 (июнь 2023) — работал, но уже с предупреждением о будущем удалении; и, наконец, ZooKeeper был физически удалён из кодовой базы в Kafka 4.0 (март 2025) — начиная с этой ветки Kafka умеет работать только в режиме KRaft. Для кластеров, которые всё ещё жили на ZooKeeper к моменту выхода 4.0, был предусмотрен путь миграции через промежуточный «мостовой» релиз (ZK-кластер переводится в KRaft постепенно, без полной остановки) — процедура и её точные шаги стоит сверять с документацией целевой версии перед миграцией, здесь важно знать сам факт, что путь без простоя существовал, а не воспроизводить устаревшие инструкции.

Живой кластер этой серии поднят на Kafka 4.3.1 — версии, где ZooKeeper не просто выключен конфигом, а физически отсутствует в образе: RemoteStorageManager-подобных остатков ZK-режима в дистрибутиве нет, поднять его на этом образе невозможно в принципе. Поэтому здесь ZooKeeper разбирается исключительно как исторический контекст, объясняющий, откуда взялся KRaft и какую проблему он решил, — не как рабочий режим, который можно включить и понаблюдать.

ZAB и Raft: родственники, но не потомки

Соблазнительно сказать «ZooKeeper использовал Raft», но это неточно и заслуживает отдельного абзаца, потому что путаница здесь встречается постоянно. Протокол консенсуса ZooKeeper называется ZAB (ZooKeeper Atomic Broadcast) — он основан на голосовании большинства и опубликован примерно в 2008 году. Raft — отдельный протокол консенсуса, опубликованный в 2014 году, на шесть лет позже. ZAB технически не может быть основан на Raft или произведён от него — он старше.

Правильная формулировка: ZAB и Raft — родственные, но не производные друг от друга протоколы. Оба решают одну и ту же задачу (реплицированный журнал с гарантией consensus при потере части узлов) и оба независимо восходят к более ранним Paxos-подобным идеям консенсуса на основе большинства голосов. Совпадение гарантий (оба требуют кворум большинства, оба переизбирают лидера при его потере, оба линеаризуют запись в лог) — следствие того, что оба протокола решают одну и ту же математическую задачу, а не то, что один списан с другого. KRaft, соответственно, не «заменил ZAB на Raft внутри той же архитектуры» — он заменил внешнюю систему консенсуса (ZooKeeper с ZAB) на встроенный Raft-кворум внутри самой Kafka, и это архитектурное решение важнее, чем формальное родство самих протоколов.

KRaft-кворум на живом кластере: почему нет split-brain

На живом трёхброкерном KRaft-кластере эксплуатационный слой Kafka виден напрямую через kafka-metadata-quorum.sh describe --status:

ClusterId:              r0mrt307TrOiv10ImOQjHw
LeaderId:               3
CurrentVoters:          [1,2,3]
CurrentObservers:       []

NodeId  Endpoint  Role      Lag
3       ...       Leader    0
1       ...       Follower  0
2       ...       Follower  0

Все три узла — voters Raft-группы __cluster_metadata, лидер (node 3) и оба фолловера синхронны (lag=0). Это тот же механизм, что и репликация партиции с данными, только применённый к самим метаданным кластера: коммит в этот лог требует голоса большинства voters, и именно из этого требования и следует ответ на вопрос «почему в Kafka нет split-brain» — уже разобранный со стороны данных в статье про репликацию. При трёх voters кворум — любые два из трёх; если сеть разбивается на две половины без большинства ни в одной из них, ни одна половина не может избрать лидера контроллера и закоммитить изменение метаданных — вместо конкурирующих «правд» о состоянии кластера получается простой до восстановления связности. Это структурное свойство Raft, а не отдельно добавленная защита поверх протокола.

flowchart TB subgraph kraft["KRaft: контроллер — Raft-группа внутри Kafka"] direction LR K1["брокер 1
+ controller (voter)"] K2["брокер 2
+ controller (voter)"] K3["брокер 3
+ controller (voter, leader)"] K1 <--> K3 K2 <--> K3 K1 <--> K2 CL["__cluster_metadata
реплицируемый лог"] K3 -.commit требует большинства.-> CL end subgraph zk["Эпоха ZooKeeper: внешний ансамбль"] direction LR B1["брокер 1"] B2["брокер 2"] B3["брокер 3"] Z1["ZK 1"] Z2["ZK 2"] Z3["ZK 3"] B1 --> Z1 B2 --> Z2 B3 --> Z3 Z1 <--> Z2 Z2 <--> Z3 end style K3 fill:#c9e4c5,stroke:#5b8a5e style K1 fill:#eee8d8,stroke:#8b7355 style K2 fill:#eee8d8,stroke:#8b7355 style CL fill:#f9f3e3,stroke:#8b7355 style B1 fill:#eee8d8,stroke:#8b7355 style B2 fill:#eee8d8,stroke:#8b7355 style B3 fill:#eee8d8,stroke:#8b7355 style Z1 fill:#f0d9c6,stroke:#b4552f style Z2 fill:#f0d9c6,stroke:#b4552f style Z3 fill:#f0d9c6,stroke:#b4552f

flowchart TB
    subgraph kraft["KRaft: контроллер — Raft-группа внутри Kafka"]
        direction LR
        K1["брокер 1
+ controller (voter)"] K2["брокер 2
+ controller (voter)"] K3["брокер 3
+ controller (voter, leader)"] K1 <--> K3 K2 <--> K3 K1 <--> K2 CL["__cluster_metadata
реплицируемый лог"] K3 -.commit требует большинства.-> CL end subgraph zk["Эпоха ZooKeeper: внешний ансамбль"] direction LR B1["брокер 1"] B2["брокер 2"] B3["брокер 3"] Z1["ZK 1"] Z2["ZK 2"] Z3["ZK 3"] B1 --> Z1 B2 --> Z2 B3 --> Z3 Z1 <--> Z2 Z2 <--> Z3 end style K3 fill:#c9e4c5,stroke:#5b8a5e style K1 fill:#eee8d8,stroke:#8b7355 style K2 fill:#eee8d8,stroke:#8b7355 style CL fill:#f9f3e3,stroke:#8b7355 style B1 fill:#eee8d8,stroke:#8b7355 style B2 fill:#eee8d8,stroke:#8b7355 style B3 fill:#eee8d8,stroke:#8b7355 style Z1 fill:#f0d9c6,stroke:#b4552f style Z2 fill:#f0d9c6,stroke:#b4552f style Z3 fill:#f0d9c6,stroke:#b4552f
KRaft: Raft-кворум контроллеров внутри Kafka — против эпохи ZooKeeper, где контроллер выбирался во внешнем ансамбле

Практическое следствие для эксплуатации — топология voters не бесплатна. В статье про репликацию разобран ровно такой случай: в минимальной трёхнодовой топологии, где брокер и контроллер живут 1:1 на одних и тех же узлах (combined), убийство двух узлов из трёх одновременно рушит и ISR партиций, и мажоритарность кворума метаданных — административные команды зависают, а не отвечают чистой ошибкой. Продакшн-кластеры обычно разносят controller-voters так, чтобы потеря нескольких брокеров данных в одной зоне отказа не задевала кворум метаданных один в один — либо через dedicated-контроллеры на отдельных узлах, либо через число voters, превышающее число брокеров в одной зоне отказа.

Метрики: consumer lag, под-репликация, latency

Consumer lag — разница между LOG-END-OFFSET (последний записанный offset партиции) и CURRENT-OFFSET (позиция, до которой группа реально закоммитила чтение) — главный симптом отставания обработки от притока данных, и первое, что стоит смотреть при любом подозрении на проблему с consumer’ом. kafka-consumer-groups.sh --describe печатает эту колонку по каждой партиции группы; сумма по всем партициям — грубый, но полезный агрегат состояния группы целиком.

Живой прогон показывает форму этой метрики во времени, а не только факт «лаг был». Продюсер пишет с темпом 20 записей/с в течение 40 секунд (800 записей всего), а consumer сперва обрабатывает первые 300 записей с искусственной задержкой 150 мс на запись (медленная фаза — имитация тяжёлой обработки), затем переключается в быстрый режим и вычитывает накопленный backlog:

Consumer lag во времени: рост при отставании, падение к 0 при догоне (Go/franz-go)0275550пик 550 (t+42с)0 (t+57с)t+8с

Суммарный лаг рос монотонно — t+8с=153, t+15с=239, t+22с=325, t+29с=421, t+36с=500 — до пика 550 на t+42с (продюсер к этому моменту уже закончил писать, а consumer ещё в медленной фазе), затем падал — t+51с=448 — и достиг ровно 0 на t+57с. Consumer в итоге обработал все 800/800 записей за 56.2 секунды. На Java (тот же сценарий, тот же кластер) картина качественно та же: пик суммарного лага 449 на t+43с, падение к 0 на t+56с, но обработано 688/688 — меньше записей от продюсера за то же окно, что объясняется различием в планировании JVM, а не ошибкой сценария.

Отдельная метрика — under-replicated partitions: партиции, у которых ISR временно меньше полного replicas. Здесь есть тонкость, которую стоит знать заранее: replica.lag.time.max.ms (дефолт 30000 мс) определяет, сколько времени фолловеру разрешено не догонять лидера, прежде чем контроллер выкинет его из ISR, — но при docker kill (SIGKILL) TCP-соединение рвётся мгновенно, а не «тихо виснет», поэтому фактическое обнаружение под-репликации на живом кластере произошло заметно раньше 30-секундного порога документации. При этом не мгновенно: первая версия скрипта проверки, ждавшая фиксированные 6 секунд перед первой проверкой, ничего не находила — under-replicated-статус контроллер выставляет не по первому же тику, поллинг пришлось растянуть до ~40 секунд, прежде чем состояние стабильно проявилось. Практический вывод: не полагайтесь на конкретное число секунд из документации как на гарантию времени реакции — под-репликация на реальном кластере наблюдалась и раньше номинального порога, и не мгновенно, а реальный интервал стоит проверять поллингом, а не однократным опросом сразу после инцидента.

Третья опора мониторинга — request latency (метрики kafka.network:type=RequestMetrics по типу запроса — Produce, FetchConsumer, FetchFollower, с перцентилями). Она не разбиралась живьём в этой серии отдельным сценарием, но именно request latency обычно первой показывает деградацию до того, как она долетает до пользовательских симптомов вроде лага: если Produce-запросы начинают заметно замедляться, это часто ранний признак либо перегрузки диска лидера, либо начавшегося ISR-shrink, который ещё не успел отразиться в списке under-replicated partitions.

Тюнинг producer: batch, linger, compression

Три параметра управляют тем, как producer формирует и отправляет батчи: batch.size (в franz-go — ProducerBatchMaxBytes) — верхняя граница размера одного батча в байтах; linger.ms — сколько миллисекунд producer готов ждать перед отправкой неполного батча в надежде накопить больше данных; compression.type — алгоритм сжатия батча целиком перед отправкой. Все три раскрывают потенциал только при асинхронной отправке — если каждая запись синхронно ждёт подтверждения предыдущей, батч вырождается в размер 1, и параметры батчинга становятся бессмысленными.

Живой прогон (топик demo-ops-tuning-producer, 5000 записей по 500 байт, четыре комбинации параметров) показывает и ожидаемый эффект, и честную аномалию:

#  batch/linger/compression   Go msg/s (MB/s)        Java msg/s (MB/s)
1  16384 / 0    / none        36973.9 (17.63)        9535.3 (4.55)
2  65536 / 20   / none        12038.4 (5.74)         10899.1 (5.20)
3  65536 / 20   / lz4         68703.4 (32.76)  MAX    10586.7 (5.05)
4  2048  / 0    / none         4543.1 (2.17)   MIN     2626.7 (1.25)  MIN

Комбинация 3 (крупный батч + linger.ms=20 + lz4) даёт максимум на обоих клиентах, а урезанный до 2048 байт батч (комбинация 4) — минимум на обоих: batch.size слишком мал, чтобы вместить сколько-нибудь заметное число 500-байтных записей с учётом служебного оверхеда записи и батча, — направление эффекта совпадает и воспроизводится. Но между комбинациями 1 и 2 клиенты расходятся: на Go комбинация 2 (linger.ms=20, без сжатия) оказалась медленнее комбинации 1 (linger.ms=0) — на коротком burst-прогоне (~2.5 МБ данных) фиксированная задержка 20 мс на каждый батч не успевает окупиться накоплением, потому что сам прогон завершается быстрее, чем набегает выгода от более крупных батчей. На Java та же комбинация оказалась быстрее baseline — «как ожидалось» по теории. Это не противоречие и не брак измерения: linger.ms — это компромисс между латентностью и throughput, а не безусловный выигрыш, и его эффект зависит от того, успевает ли реальный трафик заполнить окно ожидания достаточным объёмом данных.

func newTuningProducer(seeds []string, batchBytes int32, linger time.Duration, compression kgo.CompressionCodec, clientID string) *kgo.Client {
    cl, err := kgo.NewClient(
        kgo.SeedBrokers(seeds...),
        kgo.RequiredAcks(kgo.AllISRAcks()),
        kgo.ProducerBatchMaxBytes(batchBytes),
        kgo.ProducerLinger(linger),
        kgo.ProducerBatchCompression(compression),
        kgo.ClientID(clientID),
    )
    if err != nil {
        log.Fatalf("kgo.NewClient (tuning producer): %v", err)
    }
    return cl
}
Properties props = baseProducerProps(clientId);
props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchBytes);
props.put(ProducerConfig.LINGER_MS_CONFIG, lingerMs);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, compression);

Тюнинг consumer: fetch и max.poll.records

Со стороны consumer’а два параметра управляют размером порции данных за один вызов poll(): fetch.min.bytes — минимальный объём данных, который брокер готов накопить перед ответом на fetch-запрос (если данных меньше, брокер ждёт до fetch.max.wait.ms), и max.poll.records — верхняя граница числа записей, возвращаемых одним вызовом poll() (в franz-go — эмулируется через PollRecords(ctx, n), прямого аналога ClientOpt для этого параметра библиотека не предоставляет).

Живой прогон (топик demo-ops-tuning-consumer, 20000 записей засеяно один раз, дальше три разные consumer group читают один и тот же датасет с разными параметрами):

#  fetch.min.bytes / max.poll.records   Go avg-rec/poll   Java avg-rec/poll
1  1      / 500                         487.8             425.5
2  1      / 50                           49.9  MIN          49.8  MIN
3  65536  / 500                         487.8             416.7

max.poll.records даёт чёткий, воспроизводимый эффект на обоих клиентах — около 50 записей за вызов при значении 50, около 488 при значении 500, разница почти на порядок, ровно как и должно быть по определению параметра. А вот fetch.min.bytes эффекта не показал вовсе: 487.8 что при значении 1, что при значении 65536 байт. Честное объяснение — не «fetch.min.bytes не работает», а ограничение самой методики: топик здесь пре-засеян полностью до начала чтения, поэтому порог накопления тривиально удовлетворён уже на первом fetch-запросе независимо от значения fetch.min.bytes — данных на диске лидера и так больше любого разумного порога. Чтобы увидеть реальный эффект fetch.min.bytes, нужен сценарий с медленным непрерывным притоком данных (как в демонстрации лага выше), а не заранее заполненный топик.

cl := newGroupConsumer(seeds, topic, group,
    kgo.FetchMinBytes(fetchMinBytes),
    kgo.FetchMaxWait(fetchMaxWait),
)
// ...
fetches := cl.PollRecords(pollCtx, maxPollRecords)

Partition reassignment: ручной и автоматический — это не ребаланс группы

Здесь легко всё смешать, потому что слово «ребаланс» используется в Kafka в двух совершенно разных контекстах. Ребаланс consumer group — тема второй статьи серии: coordinator перераспределяет, какой участник группы читает какую партицию, когда состав группы меняется. Это операция над клиентами, партиции при этом физически никуда не переезжают — меняется только назначение чтения. Partition reassignment, о котором эта секция, — операция совершенно другого уровня: перемещение физических реплик партиции между брокерами. Она не имеет отношения к consumer group вообще, и производит её не coordinator, а контроллер кластера по команде оператора или автоматически.

Ручной reassignment запускается инструментом kafka-reassign-partitions.sh: оператор (или, как здесь, скрипт) строит JSON с новым размещением реплик по брокерам на основе текущего реального размещения, а не угадывает его, и передаёт этот JSON на выполнение с ограничением скорости (--throttle), чтобы перемещение данных не забило сеть и диск кластера под нагрузкой. Живой прогон: топик demo-ops-reassign (3 партиции, RF=2), заранее наполненный 2000 реальными записями — реассайнмент двигает настоящие байты, а не пустые партиции:

# JSON строится из живого describe, не угадывается заранее:
#   partition=0: [3,1] -> [1,2]
#   partition=1: [1,2] -> [2,3]
#   partition=2: [2,3] -> [3,1]
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
  --reassignment-json-file /tmp/ops-reassign.json \
  --execute --throttle 1000000

# поллинг --verify до "is completed" на всех партициях
kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
  --reassignment-json-file /tmp/ops-reassign.json --verify

Результат: 3 из 3 партиций сменили размещение реплик — describe до и после прогона подтверждает это построчно. На Java (другое стартовое размещение, тот же сценарий) — идентичный по сути факт: все 3 партиции переехали. Reassignment — операция контроллера над кластером, а не клиентской библиотеки, поэтому результат не зависит от того, каким клиентом изначально засевался топик.

flowchart LR subgraph before["ДО reassignment"] direction TB P0b["partition 0
replicas=[3,1]"] P1b["partition 1
replicas=[1,2]"] P2b["partition 2
replicas=[2,3]"] end subgraph after["ПОСЛЕ reassignment (kafka-reassign-partitions --execute)"] direction TB P0a["partition 0
replicas=[1,2]"] P1a["partition 1
replicas=[2,3]"] P2a["partition 2
replicas=[3,1]"] end before -->|"3 из 3 партиций сменили размещение"| after style P0b fill:#eee8d8,stroke:#8b7355 style P1b fill:#eee8d8,stroke:#8b7355 style P2b fill:#eee8d8,stroke:#8b7355 style P0a fill:#c9e4c5,stroke:#5b8a5e style P1a fill:#c9e4c5,stroke:#5b8a5e style P2a fill:#c9e4c5,stroke:#5b8a5e

flowchart LR
    subgraph before["ДО reassignment"]
        direction TB
        P0b["partition 0
replicas=[3,1]"] P1b["partition 1
replicas=[1,2]"] P2b["partition 2
replicas=[2,3]"] end subgraph after["ПОСЛЕ reassignment (kafka-reassign-partitions --execute)"] direction TB P0a["partition 0
replicas=[1,2]"] P1a["partition 1
replicas=[2,3]"] P2a["partition 2
replicas=[3,1]"] end before -->|"3 из 3 партиций сменили размещение"| after style P0b fill:#eee8d8,stroke:#8b7355 style P1b fill:#eee8d8,stroke:#8b7355 style P2b fill:#eee8d8,stroke:#8b7355 style P0a fill:#c9e4c5,stroke:#5b8a5e style P1a fill:#c9e4c5,stroke:#5b8a5e style P2a fill:#c9e4c5,stroke:#5b8a5e
Ручной reassignment: kafka-reassign-partitions --execute реально перемещает реплики между брокерами

Второй механизм — автоматический, и решает другую задачу: возврат лидерства к preferred-реплике (первой в списке replicas), а не перемещение самих реплик. Когда брокер с лидером партиции падает, лидерство переходит к другой реплике из ISR; когда упавший брокер восстанавливается, он снова становится просто фолловером — если auto.leader.rebalance.enable=true (дефолт), контроллер периодически (интервал — leader.imbalance.check.interval.seconds, дефолт 300 секунд) сравнивает текущих лидеров с preferred-репликами и переизбирает лидера обратно, если дисбаланс превышает порог, — без единого ручного вызова kafka-leader-election.sh.

Здесь всплывает тонкость, знакомая по статье о хранении: leader.imbalance.check.interval.secondsstatic broker config, он не поддерживает динамическое изменение через kafka-configs.sh --alter, только правку compose.yml/конфига брокера и пересоздание процесса (тот же паттерн, что log.retention.check.interval.ms в статье про retention и compaction). В живом кластере интервал статически занижен до 5 секунд специально для наблюдаемой демонстрации — с дефолтными 300 секундами (5 минут) пришлось бы просто ждать значительно дольше, чтобы увидеть эффект.

Живой прогон: preferred-реплика партиции 2 (Go-сценарий) — docker kill этого брокера → лидерство переходит на другую реплику из ISR → брокер восстанавливается (docker start) → в пределах первого же 5-секундного цикла контроллера (~3 секунды) лидерство автоматически возвращается на preferred-реплику 2. На Java — preferred 1, тот же сценарий, тот же результат: возврат за ~3 секунды.

Отдельно стоит Cruise Control — открытый инструмент, который автоматизирует и ручной reassignment, и балансировку нагрузки между брокерами по целому набору метрик (не только лидерство, но и диск, сеть, CPU), избавляя оператора от ручного построения JSON. В этой серии он разобран только обзорно: сам инструмент не разворачивался — для учебного трёхброкерного кластера полноценный Cruise Control несоразмерен задаче, а его конфигурация и логика балансировки заслуживают отдельного материала, а не пары абзацев здесь.

Rack awareness

Rack awareness заставляет контроллер размещать реплики одной партиции на разных зонах отказа (broker.rack), а не просто на разных брокерах, — цель в том, чтобы отказ целой стойки/зоны не уносил сразу несколько реплик одной партиции. Живой кластер настроен как один брокер на один rack (kafka1/rack-a, kafka2/rack-b, kafka3/rack-c), и результат подтверждает намерение: топик с 6 партициями и RF=3 получил все реплики всех 6 партиций на трёх разных rack; топик с 3 партициями и RF=2 — все реплики на двух разных rack.

Честная оговорка здесь важнее самого факта: в топологии «один брокер = один rack» связь между broker-id и rack — биекция, и любой набор из N разных broker-id автоматически покрывает N разных rack просто потому, что иначе физически не бывает. Наблюдаемый результат демонстрирует не «интеллект» алгоритма размещения, а структурное свойство конкретной топологии стенда. Полноценная проверка, где алгоритм реально выбирает между несколькими брокерами одного rack, чтобы не превысить лимит реплик на rack, требует кластера с несколькими брокерами на одном rack — за пределами объёма этого примера.

Rolling upgrade

Rolling upgrade — обновление версии Kafka по одному брокеру за раз, без общего простоя кластера: вывести узел из-под нагрузки, убедиться, что ISR остальных партиций не деградировал, обновить и перезапустить узел, дождаться, пока он догонит лог и снова войдёт в ISR, перейти к следующему. Ядро этой процедуры — контролируемый вывод и возврат одного узла в кластер под наблюдением за ISR — фактически многократно проверялось живьём как побочный эффект других сценариев этой статьи: и docker kill / docker start в демонстрации auto-leader-rebalance, и в демонстрациях под-репликации, где реплика выпадает из ISR и возвращается обратно, — во всех случаях кластер корректно детектировал уход узла, продолжал обслуживать остальные партиции и ресинхронизировал вернувшийся узел без вмешательства.

Честно: полноценный межверсионный rolling upgrade (два разных образа Kafka, постепенная замена одной версии на другую с проверкой межверсионной совместимости протокола на каждом шаге) на этом кластере не воспроизводился — используется один образ (4.3.1) на всех трёх узлах. Сама процедура для реального межверсионного апгрейда описывается текстом, а не имитируется: обновлять брокеров по одному, начиная не с текущего контроллера-лидера; после каждого шага проверять, что все партиции, где обновлённый брокер участвует, вернули полный ISR, прежде чем переходить к следующему узлу; при переходе через версии, меняющие формат внутреннего протокола или __cluster_metadata, — сверяться с матрицей совместимости конкретных версий в документации, а не полагаться на то, что «наверняка заработает».

Квоты против шумных соседей

Квоты (producer_byte_rate / consumer_byte_rate / request_percentage) ограничивают пропускную способность конкретного клиента (по client-id, пользователю или их сочетанию) — механизм защиты от ситуации, когда один producer или consumer способен насытить сеть или диск брокера и деградировать всех остальных клиентов кластера, «шумный сосед» в чистом виде.

kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type clients --entity-name quota-demo-client \
  --alter --add-config producer_byte_rate=100000

kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type clients --entity-name quota-demo-client --describe

Живой прогон вскрыл нетривиальную тонкость самого механизма квот прежде, чем дошёл до основного замера. Окно квоты (quota.window.num=11 окон по quota.window.size.seconds=1 каждое ≈ 11 секунд) при producer_byte_rate=100000 байт/с даёт около 1.1 МБ «бесплатного» burst-кредита — объём, который клиент может отправить почти без задержки, прежде чем throttling реально включится. Первая попытка замера (2 МБ) throttling не показала вовсе — ни на franz-go, ни на официальном kafka-producer-perf-test.sh — объём оказался слишком близок к порогу burst-кредита, чтобы эффект был однозначно виден. Финальный прогон выбран заведомо больше порога: 8000 записей по 1000 байт = 8 МБ, топик demo-ops-quota:

producer_byte_rate quota: throughput до и после (характерный прогон, шкала логарифмическая)10100100010000100000 msg/s16712.69052.6327.0322.7Go, БЕЗ quotaJava, БЕЗ quotaGo, С quotaJava, С quota

Без квоты: Go — 16712.6 сообщений/с (15.94 МБ/с), Java — 9052.6 сообщений/с (8.63 МБ/с). С квотой producer_byte_rate=100000 (~97.7 КБ/с): Go — 327.0 сообщений/с (0.31 МБ/с), Java — 322.7 сообщений/с (0.31 МБ/с) — падение примерно в 51 раз на Go и 28 раз на Java. Показательнее самого падения — то, что throttled-throughput у обоих клиентов практически совпал (0.31 МБ/с у обоих), хотя без квоты их «естественный» throughput различался почти вдвое: квота — ограничение на стороне брокера, она не зависит от клиентской библиотеки и одинаково режет любого, кто пытается превысить назначенный порог. Итоговые ~0.31 МБ/с немного выше номинальных ~0.0977 МБ/с — остаточный эффект того самого burst-кредита окна квоты, действующий и внутри самого прогона, не только на старте.

Продакшн-чеклист

  • RF и min.insync.replicas согласованы с ожидаемым числом одновременных отказов — RF=3/min.insync.replicas=2 переживает отказ одного брокера без потери availability; подробности гарантий — в статье про репликацию.
  • Consumer lag мониторится по группам и алертится на рост, а не только на абсолютное значение — лаг, который растёт быстрее, чем обрабатывается, симптом раньше, чем лаг, который просто велик, но стабилен.
  • Under-replicated / offline partitions — алерт на любое ненулевое значение, не дожидаясь replica.lag.time.max.ms: как показано выше, реальное обнаружение под-репликации может произойти и раньше, и позже наивного ожидания по документации.
  • KRaft-кворум контроллеров мониторится отдельно от ISR партиций данныхkafka-metadata-quorum.sh describe --status, число voters и их lag; при combined-топологии (controller+broker на одних узлах) явно понимать, что потеря большинства узлов рушит и данные, и метаданные одновременно.
  • Static-конфиги задокументированы отдельно от dynamicleader.imbalance.check.interval.seconds, log.retention.check.interval.ms и подобные требуют правки конфигурации брокера и его пересоздания, а не kafka-configs.sh --alter; неверное ожидание здесь стоит часов отладки «почему конфиг не применился».
  • Reassignment — всегда с --throttle на проде, чтобы перемещение данных не забило сеть/диск под живой нагрузкой; и всегда JSON, построенный из текущего реального размещения, а не из предположений.
  • Rack awareness подтверждена не на бумаге, а --describe — реальное распределение реплик по broker.rack, а не предположение, что оно настроено правильно.
  • Квоты назначены заранее для известных «тяжёлых» клиентов (batch-джобы, миграции, ETL), а не постфактум после инцидента с одним шумным соседом.
  • Backup конфигурации топиков (kafka-configs.sh --describe по всем топикам, версионируется отдельно от самих данных) — конфиг топика не реплицируется вместе с данными на уровне «просто скопировать диск».

Собранная воедино картина: KRaft заменил внешний ZooKeeper-кворум на встроенный Raft-кворум контроллеров, и это же структурное свойство Raft — требование голоса большинства — одновременно объясняет и отсутствие split-brain, и то, почему combined-топология уязвима к потере большинства узлов сразу с двух сторон. Метрики (lag, under-replicated, latency), тюнинг producer/consumer и квоты — это то, чем управляют изо дня в день; reassignment, rack awareness и rolling upgrade — то, что происходит реже, но требует понимания, что именно перемещается и кем.

Дальше в серии — как Kafka интегрируется с внешним миром без написания собственного connector-кода и как схемы данных эволюционируют без разрыва совместимости, в статье про экосистему: Kafka Connect и Schema Registry. Тонкости самой JVM под нагрузкой продюсера/консьюмера, GC-паузы и их влияние на латентность — на границе с этой темой, в статье «JVM и messaging: Kafka». Гео-репликация и disaster recovery поверх модели ISR и KRaft-кворума, разобранных здесь и в статье про репликацию, — в финальной статье серии про MirrorMaker 2 и геораспределённый Kafka. Для навигации по всей теме messaging — карта messaging-landscape-map.

Обсуждение в Telegram

Присоединиться →

Комментарии