Репликация и надёжность Kafka: ISR, acks, durability

Как Kafka не теряет сообщения: лидеры и фолловеры, ISR, acks (0/1/all), min.insync.replicas, unclean leader election и идемпотентный producer — и как одной неверной настройкой превратить надёжный кластер в дырявый

Надёжность Kafka — не свойство «из коробки», а результат согласованных настроек с двух сторон. Можно поднять кластер с фактором репликации 3 и всё равно терять данные, если producer шлёт с acks=1, а брокер допускает unclean leader election. Durability в Kafka живёт на пересечении acks, min.insync.replicas и репликации — и понимать это пересечение важнее, чем знать любую отдельную настройку.

Это третья статья серии «Kafka: глубокое погружение». Опирается на модель партиций из первой статьи и на модель чтения из статьи про consumer groups — здесь речь о том, что происходит с записью до того, как она вообще становится видна consumer’у. Все числа ниже — из живого трёхброкерного KRaft-кластера (docker kill реальных контейнеров, не симуляция), включая три честных попытки, которые не удалось довести до воспроизведения. Весь код этого разбора открыт в репозитории digital-cookbook — стенды kafka/go/replication (franz-go) и kafka/java/replication (kafka-clients), плюс оркестрирующий скрипт ops/broker-kill.sh.

Репликация и ISR: падение брокера, выборы нового лидера и ноль потерь при acks=all

В статье

Лидер и фолловеры: репликация партиции

У каждой партиции — ровно один лидер (брокер, который принимает все записи и все чтения этой партиции в данный момент) и ноль или больше фолловеров — реплик, которые непрерывно вытягивают (fetch) записи с лидера и дописывают их в свой локальный лог в том же порядке. Клиент никогда не пишет напрямую фолловеру и, по умолчанию, не читает с него: фолловер существует ради одной цели — быть готовым стать лидером, если текущий лидер откажет.

Replication factor (RF) — число копий партиции, включая лидера. RF=3 означает: лидер плюс два фолловера, размещённые на разных брокерах контроллером (с учётом rack awareness, если она настроена — подробнее в статье про эксплуатацию и KRaft). RF — это про то, сколько копий существует физически; сколько из них обязаны реально подтвердить запись, прежде чем producer получит успех, — отдельная настройка, acks, и путать эти две оси одна из самых частых ошибок при проектировании надёжности.

ISR — кто считается «в синхроне»

Не все реплики партиции равноценны в любой момент времени. ISR (in-sync replicas) — подмножество реплик (включая лидера), которые фактически догнали лидера в пределах допустимого лага, определяемого replica.lag.time.max.ms (сколько времени фолловеру разрешено не присылать fetch-запрос или не догонять лидера, прежде чем контроллер сочтёт его отставшим и выкинет из ISR). Фолловер, который отстал сильнее — из-за медленного диска, перегруженной сети или просто временного падения, — выпадает из ISR, но продолжает существовать как реплика: он остаётся в списке replicas, просто временно не считается «достаточно свежим», чтобы участвовать в гарантиях durability.

Именно размер ISR, а не размер RF, определяет реальную устойчивость к потере данных в моменте. Партиция с RF=3, но ISR, схлопнувшимся до одного элемента (лидера самого с собой) из-за упавших фолловеров, физически имеет только одну копию данных — RF=3 на бумаге ничего не гарантирует, если ISR не отражает эту цифру прямо сейчас. Контроллер (в KRaft — через Raft-группу __cluster_metadata, подробности — дальше в серии) — единственная сторона, которая может официально сжать или расширить ISR: сам брокер-лидер не решает это в одиночку, он лишь сообщает контроллеру наблюдаемое состояние фолловеров через AlterPartition.

acks: что на самом деле подтверждает producer

acks — параметр producer’а, а не брокера: он определяет, дожидается ли клиент подтверждения записи и от кого именно, прежде чем считать её отправленной.

  • acks=0 — producer не ждёт вообще ничего: запись уходит на лидера, и клиент немедленно продолжает, даже не зная, дошла ли она до брокера физически. Максимальная скорость, нулевая гарантия — включая гарантию того, что запись вообще была принята.
  • acks=1 — producer ждёт подтверждения от лидера, что запись физически дописана в его локальный лог. Фолловеры при этом ещё не обязаны были её забрать: если лидер откажет ровно в промежутке между «записал у себя» и «фолловер успел зафетчить», подтверждённая клиенту запись пропадает вместе с лидером.
  • acks=all (синоним acks=-1) — producer ждёт подтверждения от лидера, что запись реплицирована на все текущие ISR, а не просто на какое-то фиксированное число реплик. Это отличие важно: acks=all — гарантия относительно ISR в момент записи, а не относительно исходного RF, и если ISR уже сжат до одного элемента, acks=all формально выполняется на этой единственной копии.

Разница в порядке величины видна прямо на живом стенде. Топик demo-repl (1 партиция — намеренно, чтобы лидер был однозначным, RF=3, min.insync.replicas=2), 200 сообщений на уровень, последовательный ProduceSync (каждая запись ждёт ответа на предыдущую, прежде чем уйти):

acks=0 vs acks=1 vs acks=all: throughput (Go/franz-go, характерный прогон, шкала логарифмическая)100 msg/s1000 msg/s10000 msg/s6715.3 msg/s1463.8 msg/s364.8 msg/sacks=0acks=1acks=all

acks=0 дал 6715.3 сообщений в секунду, acks=1 — 1463.8, acks=all — 364.8: почти 18-кратный разрыв между крайними значениями, монотонно в ожидаемую сторону (0 не ждёт ничего → 1 ждёт диск лидера → all ждёт весь ISR). Абсолютные числа host-зависимы (один docker-хост, без реальной сетевой задержки между «разными» брокерами) — важен не факт «6715», а порядок величины и направление эффекта, которое воспроизведётся на любом кластере.

На Java (kafka-clients, тот же сценарий) абсолютные числа мельче и порядок не такой чистый: acks=0 — elapsed=1.390с, throughput=143.9 msg/s; acks=1 — elapsed=1.279с, throughput=156.3 msg/s; acks=all — elapsed=1.531с, throughput=130.6 msg/s. Здесь есть честная оговорка: acks=0 в этом прогоне — не самый быстрый уровень, хотя протокольно обязан быть таковым. Причина — методологическая, не протокольная: сценарий тестирует уровни строго по порядку 0 → 1 → all в одном процессе, и acks=0 идёт первым, до прогрева JIT JVM. Содержательный факт, который порядок Java-чисел действительно подтверждает, — acks=all стабильно медленнее всего в обоих клиентах, а вот сравнение 0 и 1 между собой на Java смазано прогревом и не может служить доказательством разницы именно этих двух уровней.

func newProducer(seeds []string, acks kgo.Acks, idempotent bool, recordRetries int, reqTimeout time.Duration) *kgo.Client {
    opts := []kgo.Opt{
        kgo.SeedBrokers(seeds...),
        kgo.RequiredAcks(acks),
        kgo.ProduceRequestTimeout(reqTimeout),
    }
    if !idempotent {
        opts = append(opts, kgo.DisableIdempotentWrite())
    }
    // ...
}

func acksFromString(s string) kgo.Acks {
    switch s {
    case "0":
        return kgo.NoAck()
    case "1":
        return kgo.LeaderAck()
    case "all", "-1":
        return kgo.AllISRAcks()
    }
    // ...
}
private static Properties producerProps(String acks, boolean idempotent, int retries, int reqTimeoutMs) {
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Kafka.BOOTSTRAP);
    props.put(ProducerConfig.ACKS_CONFIG, acks); // "0" | "1" | "all"
    props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, idempotent);
    props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, reqTimeoutMs);
    if (retries >= 0) {
        props.put(ProducerConfig.RETRIES_CONFIG, retries);
    }
    return props;
}

Вот как выглядит путь одной записи через ISR при acks=all — брокер не отвечает producer’у, пока не наберёт подтверждений от всех реплик, входящих в ISR на момент записи:

flowchart LR P["Producer
acks=all"] -->|"1. запись"| L["Лидер партиции"] L -->|"2. append в локальный лог"| LD[("диск лидера")] L -.->|"3. fetch"| F1["Фолловер 1
(в ISR)"] L -.->|"3. fetch"| F2["Фолловер 2
(в ISR)"] F1 -->|"4. подтвердил offset"| L F2 -->|"4. подтвердил offset"| L L -->|"5. весь ISR подтвердил → ack"| P style P fill:#f9f3e3,stroke:#8b7355 style L fill:#c9e4c5,stroke:#5b8a5e style F1 fill:#eee8d8,stroke:#8b7355 style F2 fill:#eee8d8,stroke:#8b7355

flowchart LR
    P["Producer
acks=all"] -->|"1. запись"| L["Лидер партиции"] L -->|"2. append в локальный лог"| LD[("диск лидера")] L -.->|"3. fetch"| F1["Фолловер 1
(в ISR)"] L -.->|"3. fetch"| F2["Фолловер 2
(в ISR)"] F1 -->|"4. подтвердил offset"| L F2 -->|"4. подтвердил offset"| L L -->|"5. весь ISR подтвердил → ack"| P style P fill:#f9f3e3,stroke:#8b7355 style L fill:#c9e4c5,stroke:#5b8a5e style F1 fill:#eee8d8,stroke:#8b7355 style F2 fill:#eee8d8,stroke:#8b7355
acks=all: producer получает ответ только после того, как ВСЕ реплики из текущего ISR подтвердили запись

min.insync.replicas и живой failover: 30 записей переживают docker kill

min.insync.replicas — конфиг топика, а не producer’а: минимальный размер ISR, при котором брокер вообще принимает запись с acks=all. Если текущий ISR меньше этого порога, запись отклоняется явной ошибкой NOT_ENOUGH_REPLICAS (или NOT_ENOUGH_REPLICAS_AFTER_APPEND, если ISR сжался уже после того, как лидер физически дописал запись) — брокер осознанно отказывает в записи, вместо того чтобы молча принять её на недостаточном числе копий. Комбинация, которая реально не теряет подтверждённые записи при отказе одного брокера, — RF=3, acks=all, min.insync.replicas=2: если один из трёх узлов пропадает, оставшиеся два всё ещё формируют ISR не меньше порога, и запись либо проходит с полной гарантией, либо явно отклоняется — но никогда не проходит «наполовину».

Это ядро стенда — не гипотетическая проверка, а реальный docker kill живого лидера под нагрузкой. Топик demo-repl-durability (RF=3, min.insync.replicas=2), 30 сообщений с acks=all.

scenario_durability() {
  local topic=demo-repl-durability
  client_run -scenario=setup -topic="$topic" -partitions=1 -rf=3 -minisr=2

  echo "--- produce 30 сообщений acks=all ДО падения ---"
  client_run -scenario=produce -topic="$topic" -n=30 -acks=all -idempotent=true -prefix=durable

  local leader_id
  leader_id=$(client_run -scenario=describe -topic="$topic" | grep -oE 'leader=[0-9]+' | grep -oE '[0-9]+')

  # docker kill (SIGKILL), НЕ docker stop: controlled shutdown сам переносит
  # лидерство ДО завершения процесса — это "красивое" отключение, не крах.
  echo "--- docker kill kafka-cookbook-$leader_id (реальный крах без graceful shutdown) ---"
  docker kill "kafka-cookbook-$leader_id"

  echo "--- verify: count после failover должен == 30 ---"
  client_run -scenario=verify -topic="$topic" -expect=30
}

На Go: acked=30/30 до отключения → docker kill контейнера с лидером (kafka-cookbook-3, node.id=3) → реальный failover, лидер сменился 3 → 1, ISR сжался [3,1,2] → [1,2], узел 3 помечен offline → verify читает ровно 30/30, ни одной потери. На Java (отдельный прогон, demo-repl-java-durable) — идентичный результат с другими node.id: failover лидера 1 → 2, ISR [1,2,3] → [2,3], verify OK: 30/30. Тот же исход на обоих клиентах — это не совпадение, а прямое следствие того, что failover и ISR — свойства брокера и контроллера, а не клиентской библиотеки.

sequenceDiagram participant Prod as producer (acks=all) participant L as лидер (node 3) participant F1 as фолловер (node 1) participant F2 as фолловер (node 2) participant C as контроллер (Raft-кворум) Prod->>L: 30 записей acks=all L->>F1: репликация L->>F2: репликация L-->>Prod: acked=30/30 Note over L: docker kill (SIGKILL) — без предупреждения кластеру C->>C: обнаруживает пропажу лидера (heartbeat истёк) C->>F1: избран новый лидер (был в ISR) Note over C: ISR [3,1,2] -> [1,2], node 3 offline Prod->>F1: verify: чтение с начала F1-->>Prod: 30/30 — ни одна acked-запись не потеряна

sequenceDiagram
    participant Prod as producer (acks=all)
    participant L as лидер (node 3)
    participant F1 as фолловер (node 1)
    participant F2 as фолловер (node 2)
    participant C as контроллер (Raft-кворум)

    Prod->>L: 30 записей acks=all
    L->>F1: репликация
    L->>F2: репликация
    L-->>Prod: acked=30/30

    Note over L: docker kill (SIGKILL) — без предупреждения кластеру
    C->>C: обнаруживает пропажу лидера (heartbeat истёк)
    C->>F1: избран новый лидер (был в ISR)
    Note over C: ISR [3,1,2] -> [1,2], node 3 offline

    Prod->>F1: verify: чтение с начала
    F1-->>Prod: 30/30 — ни одна acked-запись не потеряна
Живой broker-kill: docker kill лидера → детекция контроллером → failover → verify без потерь

Отдельно проверена и обратная сторона verify — фикс, сделавший checkdup жёстким: раньше расхождение только печаталось, теперь при дублях процесс падает (exit 1) с точным списком дублирующихся индексов, а не просто предупреждением, которое легко пропустить в логе. Проверено вживую в обе стороны: зелёный путь (10 уникальных индексов, exit 0) и путь падения (искусственно продублированные 5 индексов записаны дважды каждый — 10 записей, 5 дублей, verify -checkdup завершается exit 1 со списком [0 1 2 3 4]).

Идемпотентный producer: retries без дублей

retries (и таймаут доставки delivery.timeout.ms вокруг них) решают проблему потери записи при временном сбое — но создают новую: если producer не получил ответ на попытку (таймаут, обрыв соединения) и ретраит, а на самом деле первая попытка дошла, но ответ потерялся на обратном пути, повторная отправка создаёт дубль в логе. Без дополнительного механизма это неизбежное следствие retries поверх ненадёжной сети — не баг конкретного клиента.

enable.idempotence=true закрывает именно эту дыру, а не гарантии acks в целом: каждому producer’у назначается PID (producer ID) и sequence-номер на партицию, брокер отслеживает последний принятый sequence и отбрасывает повторы с уже виденным номером как no-op вместо повторной записи. Идемпотентность работает строго в пределах одной сессии producer’а на одну партицию — она не решает дедупликацию между разными producer-инстансами и не заменяет транзакции (мостик к exactly-once semantics, где идемпотентный producer — фундамент, а не вся картина). Технически идемпотентность требует acks=all: и Go, и Java клиент в стенде явно отключают её, если тестируется acks=0 или acks=1, потому что без всех ISR-подтверждений sequence-трекинг не может гарантировать то, что обещает.

Проверить дубль под реальным failover’ом — честная попытка, а не гипотеза. Топик demo-repl-idem-on (enable.idempotence=true, acks=all), 500 записей с задержкой 5 мс между ними, docker kill лидера ровно посередине сценария: итог acked=499 failed=1 (единственная ошибка — COORDINATOR_LOAD_IN_PROGRESS на холодном старте, не связана с самим kill), и ноль дублей среди 499 уникальных записей. Контрольный прогон с enable.idempotence=false на том же сценарии (demo-repl-idem-off): acked=500/500, один запрос буквально завис на границе kill на 9.469265944 секунды (прямое доказательство того, что retry/failover-механизм действительно сработал под нагрузкой) — и всё равно ноль дублей среди 500 записей.

Честно: дубль здесь не воспроизведён ни разу — ни при выключенной идемпотентности, ни, тем более, при включённой. Причина не в том, что стенд не пытался, а в самой природе гонки: дубль возникает, когда ack от брокера физически теряется после того, как запись уже закоммичена на диск, — эта конкретная узкая гонка (не «брокер не ответил вовсе», а «брокер ответил, но ответ не долетел») не гарантированно ловится внешним docker kill, который обрывает соединение, а не выборочно теряет один TCP-пакет с ответом. Ретрай-механизм в обоих прогонах реально сработал (это видно по зависшему на 9.5 секунды запросу) — просто ни разу не сработал именно в том узком окне, которое создало бы дубль без идемпотентности.

acks=1: окно потери, которое не удалось поймать

Теоретическое окно потери у acks=1 понятно из самого протокола: producer получает подтверждение сразу после того, как лидер физически дописал запись в свой лог, — до того, как хотя бы один фолловер успел её зафетчить. Если лидер отказывает строго в этом промежутке, подтверждённая клиенту запись существует только на упавшем диске и пропадает вместе с ним; новый лидер избирается из фолловеров, которые её никогда не видели. Гарантия отсутствует не потому, что что-то сломано, а потому что acks=1 по определению не спрашивает фолловеров вовсе.

Поймать эту потерю живьём стенд пытался трижды, нарастающей агрессивностью — и ни разу не поймал:

  1. Синхронный producer (RecordRetries(0), acks=1), kill через 1–1.5 секунды после старта, до 3000 записей — 0 потерь.
  2. Асинхронный fire-and-forget без ограничения скорости, до 500 000 записей, kill через 0.3–0.4 секунды — 0 потерь: пачка успевала физически подтвердиться раньше, чем срабатывал kill (пайплайнинг franz-go считает счёт на десятки-сотни тысяч сообщений в секунду, окно между «отправлено» и «подтверждено» слишком узкое, чтобы docker успел его накрыть).
  3. Throttled fire-and-forget (искусственная задержка 5 мс между отправками, чтобы растянуть окно), топик demo-repl-ff4, 1000 записей, kill через 2.5 секунды: callback-успехов=176/1000, callback-ошибок=84, без ответа=740. После рестарта consumer прочитал ровно те же 176 индексов (0–175), что и были подтверждены callback’ом, — ни одна acked-запись не потеряна.
// produceFireAndForget шлёт n записей АСИНХРОННО и даёт коллбэкам collectFor
// времени долететь. Именно так выглядит наибольший риск потери
// неотреплицированного хвоста при падении лидера: пока producer шлёт пачку
// без бэкпрешера, несколько записей могут быть физически на диске лидера и
// уже подтверждены клиенту, но ещё не прочитаны фетчем фолловера.
func produceFireAndForget(cl *kgo.Client, topic string, n int, valuePrefix string, collectFor, delay time.Duration) (acked, failed []int) {
    for i := 0; i < n; i++ {
        if delay > 0 && i > 0 {
            time.Sleep(delay)
        }
        rec := &kgo.Record{Topic: topic, Value: []byte(fmt.Sprintf("%s-%d", valuePrefix, i))}
        cl.Produce(context.Background(), rec, func(r *kgo.Record, err error) {
            // в реальном коде — под mutex дописать i в acked или failed;
            // ack долетает до callback ДО или ПОСЛЕ kill — вот и вся гонка
        })
    }
    time.Sleep(collectFor)
    return acked, failed
}

Вывод — не оправдание, а честная фиксация: окно потери у acks=1 существует по дизайну протокола, но на быстрой локальной docker-сети с частым fetch-поллингом фолловеров это окно оказалось у́же, чем удалось поймать docker kill «снаружи». Гарантии как не было, так и нет — это не значит, что acks=1 безопасен на практике, это значит, что для демонстрации его небезопасности нужна более прицельная гонка (например, искусственная задержка фетча фолловера), чем внешний обрыв контейнера.

Асимметрия здесь честная и однонаправленная: сценарий produce-ff (специально придуманный именно для этой охоты) реализован только на Go. Java-эквивалент сознательно не строился — окно узкое независимо от языка клиента, и повтор того же эксперимента на другом клиенте дал бы тот же отрицательный результат, не добавив нового знания. Все остальные сценарии этого стенда (throughput, durability-failover, min.insync.replicas, идемпотентность) прогнаны на обоих клиентах.

Unclean leader election: durability vs availability

Если все реплики из ISR становятся недоступны одновременно (например, все три брокера партиции разом упали), у контроллера есть развилка: ждать, пока хотя бы одна из них вернётся (партиция недоступна для записи всё это время — приоритет durability), либо выбрать лидера из реплики, которая когда-то выпала из ISR и заведомо отстала (партиция снова доступна немедленно, но потенциально теряет все записи, которых у этой отставшей реплики никогда не было, — приоритет availability). Второй путь и называется unclean leader election, и он выключен по умолчанию (unclean.leader.election.enable=false) ровно потому, что тихая потеря данных — куда более коварная авария, чем видимый простой: простой заметен сразу, а тихая потеря обнаруживается только тогда, когда кто-то замечает, что данных не хватает.

Механизм переключения проверен на живом кластере: unclean.leader.election.enable=true, выставленный через AlterConfigs на уровне топика, действительно применяется динамически поверх кластерного дефолта false — это подтверждает kafka-configs.sh --describe сразу после смены.

kafka-configs.sh --bootstrap-server kafka1:9092 --entity-type topics \
  --entity-name demo-repl --alter \
  --add-config unclean.leader.election.enable=true

kafka-configs.sh --bootstrap-server kafka1:9092 --entity-type topics \
  --entity-name demo-repl --describe
# unclean.leader.election.enable=true (DYNAMIC_TOPIC_CONFIG) поверх
# кластерного дефолта false

Сам сценарий потери данных при unclean-выборе — честно не продемонстрирован, и причина не техническая случайность, а прямое следствие топологии этого стенда. Чтобы вызвать unclean election, нужно убить все актуальные ISR-реплики партиции, оставив в живых только отставшего фолловера, — а в трёхнодовой топологии (RF=3, три брокера всего) это означает убить как минимум две ноды из трёх. Именно это действие — «убить 2 из 3» — ломает не только ISR партиции, но и кворум контроллера целиком: в этой конфигурации контроллер физически не может провести никакое переизбрание лидера, ни чистое, ни грязное, потому что ему для этого самого нужен кворум голосов. Разбор этого — уже не про репликацию партиции, а про кворум метаданных кластера — дальше в статье и полностью развёрнут в статье про KRaft и эксплуатацию.

Split-brain и кворум контроллера: убить 2 из 3

Это самая содержательная находка стенда: попытка буквально воспроизвести сценарий «убить 2 из 3 реплик при min.insync.replicas» дала результат, который не совпал с ожиданием по документации — и разбор несовпадения важнее, чем сам факт совпадения ожиданий.

Сценарий 4a — буквально по теории протокола. Топик demo-repl-minisr-literal, RF=3, min.insync.replicas=2, убиты 2 брокера из 3. Ожидание по теории протокола: ISR схлопывается до 1, acks=all получает чистый NOT_ENOUGH_REPLICAS. Реальность другая: producer не получил чистую ошибку вообще — вместо неё record failed after being retried too many times за elapsed=6.070149501s, то есть исчерпание ретраев по таймауту, а не осмысленный код ошибки протокола.

Причина — в топологии самого кластера, не в баге клиента. Все три ноды стенда — combined broker+controller, и все три одновременно являются voters Raft-кворума метаданных (__cluster_metadata). Убийство 2 из 3 нод ломает не только ISR данных, но и мажоритарность кворума контроллера тоже: остаётся 1 живой voter из 3, а Raft требует голоса большинства, чтобы закоммитить что угодно в metadata log — включая AlterPartition, которым брокер-лидер сообщает контроллеру «фолловер отстал, сожми ISR». Без большинства это изменение физически не может быть подтверждено: describe после kill показывает устаревшее состояние (ISR не сжат, offline=[]), потому что контроллер не в состоянии зафиксировать новую правду, даже если она ему технически известна.

Проверено и дополнительно: сами административные команды kafka-metadata-quorum.sh describe --status и kafka-topics.sh --describe при живых 2 из 3 voters зависают целиком (DisconnectException: node 3 disconnected) — это не частный случай продюсера, а системный эффект потери кворума на весь путь метаданных кластера.

scenario_minisr_literal() {
  local topic=demo-repl-minisr-literal
  client_run -scenario=setup -topic="$topic" -partitions=1 -rf=3 -minisr=2

  echo "--- убиваю 2 брокера из 3 (оставляю живым только один) ---"
  docker kill kafka-cookbook-2 kafka-cookbook-3
  sleep 3

  # ISR НЕ подтверждён контроллером — останется устаревшим (кворум потерян)
  client_run -scenario=describe -topic="$topic" || true

  # ожидаем таймаут/исчерпание ретраев, НЕ чистый NOT_ENOUGH_REPLICAS
  client_run -scenario=minisr-produce -topic="$topic" -acks=all || true
}

Сценарий 4b — контраст, где теория подтверждается буквально. Топик demo-repl-minisr-contrast, RF=2 (не 3!), min.insync.replicas=2, убита ровно одна реплика — кворум контроллера цел (живы 2 брокера из 3, большинство есть). В этой топологии ISR корректно сжимается контроллером до [1], и producer с acks=all получает именно чистый NOT_ENOUGH_REPLICAS — за 58 мс на Go (elapsed=58.074868ms) и за 813 мс на Java (elapsed=812.9ms; при повторном прогоне после фикса стенда — 679.1 мс). Ровно как в теории протокола — потому что здесь ISR-подсистема и кворум контроллера не связаны одним и тем же отказом.

4a (RF=3, kill 2 из 3, кворум контроллера потерян):
  [minisr-produce] topic=demo-repl-minisr-literal acks=all ->
    record failed after being retried too many times (elapsed=6.070149501s)

4b (RF=2, kill 1 из 2 реплик, кворум контроллера цел):
  [minisr-produce] topic=demo-repl-minisr-contrast acks=all ->
    NOT_ENOUGH_REPLICAS (elapsed=58.074868ms)   # Go
    NOT_ENOUGH_REPLICAS (elapsed=812.9ms)       # Java

Итоговый урок шире, чем частность конкретного стенда: min.insync.replicas защищает от потери данных только если контроллер способен обработать ISR-изменение. В минимальной связной трёхнодовой топологии, где брокер и контроллер живут 1:1 на одних и тех же нодах, потеря большинства узлов ломает обе подсистемы разом — availability деградирует раньше и грязнее, чем «просто получите NOT_ENOUGH_REPLICAS». Это не недостаток KRaft, а прямая причина, по которой продакшн-кластеры обычно разносят controller-voters так, чтобы потеря нескольких брокеров с данными не задевала кворум метаданных один в один (dedicated controllers либо больше voters, чем брокеров в одной зоне отказа) — тема, разворачиваемая в статье про эксплуатацию и тюнинг KRaft.

Здесь же лежит ответ на вопрос «почему в Kafka нет split-brain», который стоит явно отделить от ISR-репликации данных: split-brain на уровне метаданных кластера предотвращает не репликация партиций, а Raft-кворум контроллера. Контроллер — это отдельная Raft-группа (__cluster_metadata), и коммит в неё требует голоса большинства voters. Именно поэтому при потере связности сеть без большинства не может избрать лидера контроллера ни в одной из своих половин: вместо двух конкурирующих «правд» о состоянии кластера получается простой до восстановления связности — структурное свойство Raft, а не отдельная защита, добавленная поверх. Практическое наблюдение этого механизма на живом трёхвотерном кворуме, включая устойчивость к отказу одного voter из трёх без потери availability, — в статье про KRaft, метрики и тюнинг, которая логически продолжает эту находку.


Собранная воедино картина: acks определяет, чего дожидается producer, min.insync.replicas определяет, когда брокер вообще соглашается принять запись, а ISR — это подвижная граница между «формально реплика существует» и «реплика реально готова подтвердить durability прямо сейчас». Все три оси совместно управляемы — и все три ломаются одинаково грубо, если забыть, что кворум метаданных кластера и репликация данных, хоть и выглядят как разные темы, в минимальной топологии физически завязаны на одни и те же узлы.

Дальше в серии: как партиция ведёт себя на диске — сегменты, retention и log compaction — в статье про хранение и retention. Идемпотентный producer, разобранный здесь, — фундамент для exactly-once semantics и транзакций. Эксплуатационная сторона кворума контроллера — метрики, reassignment, история KRaft — в статье про эксплуатацию и тюнинг, а гео-репликация поверх ISR-модели, разобранной здесь (stretch-кластер и MirrorMaker 2), — в финальной статье серии про гео-репликацию. Прикладной угол на надёжность сообщений — не только внутри Kafka, но и на стыке с брокером и БД, — в статье «Транзакции и брокеры: RabbitMQ и Kafka». Для навигации по всей теме messaging — карта messaging-landscape-map.

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

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

Комментарии