В Kafka сообщение считается закоммиченным не тогда, когда его записал лидер, а тогда, когда его успели забрать достаточно реплик. Границей этой «достаточности» служит high watermarkМаксимальное смещение, до которого данные считаются закоммиченными, потому что их успели забрать достаточно реплик., а множеством реплик, которые обязаны успеть, — ISRНабор реплик, которые лидер считает синхронными себе; сообщение считается закоммиченным, когда его успели забрать все реплики из этого набора.. Вопрос, который разбираем: как именно лидер решает, кого держать в ISR, как вычисляется HW и что происходит, когда брокер выпадает из строя или возвращается с отставшим логом. Нетривиальность в том, что все эти решения принимаются асинхронно и распределённо, а цена ошибки — либо потеря подтверждённых данных, либо вечная остановка реплики, которая никак не может догнать растущий HW.
Кто попадает в ISR и как из него выпадают
ISR. Лидер проверяет три условия по порядку, с коротким замыканием: по реплике не идёт незавершённая операция изменения ISR (состояние inflightоперация изменения ISR уже запущена и ещё не завершена), её идентификатора ещё нет в текущем ISR, и она проходит проверку isReplicaIsrEligible. Порядок такой, потому что более дешёвые проверки состояния отсекают реплику раньше, чем выполняется обращение к метаданным брокера.
isReplicaIsrEligible определяет, имеет ли реплика право войти в ISR. Сначала она достаёт реплику из состояния лидера; если её там нет (например, топик уже удалён), реплика в ISR не допускается. Затем проверяются три вещи: реплика не fencedпомечена как недоступная для лидерства, не находится в controlled shutdownуправляемом завершении работы, и её сохранённая эпоха брокераbroker epoch — номер регистрации брокера в контроллере совпадает с кэшированной в метаданных. Совпадение эпох проверяется так:
storedBrokerEpoch.isPresent && cachedBrokerEpoch.isPresent &&
(storedBrokerEpoch.get == -1 || storedBrokerEpoch.get == cachedBrokerEpoch.get)Значение -1 — это явный обход проверки эпохи: и в проверке допуска, и на стороне контроллера оно означает «эпоху не сверяем».
Само понятие синхронности задаёт isFollowerInSync — функция, которая по состоянию реплики-фолловера определяет, догнала ли она лидера, и используется при решении о добавлении реплики в ISR. У лидера должен быть локальный лог, а последний зафиксированный конец лога реплики (followerEndOffset) должен быть не меньше HW лидера и не меньше начала текущей эпохи лидера. Порог задан именно так, а не числовым лагом, потому что реплика обязана догнать committed-границу и текущую эпоху: иначе, если между HW и LEO лидера есть закоммиченные данные, реплика может стать лидером до того, как их получит, и данные потеряются. Технически реплика не должна быть в ISR при отставании дольше replicaLagTimeMaxMs даже при LEO ≥ HW, но ради согласованности с тем, как синхронность определяет сам фолловер, проверяется только HW.
Выпадение из ISR устроено иначе. needsShrinkIsr вызывает getOutOfSyncReplicas(replicaLagTimeMaxMs) — то есть порогом maxLagMs служит replicaLagTimeMaxMs. Для каждой реплики-кандидата (все из ISR, кроме локального брокера) isFollowerOutOfSync считает реплику отставшей, если её нет в состоянии лидера, иначе — если она не догнала лидера по LEO, текущему времени и порогу. Если операция изменения ISR уже в полёте, getOutOfSyncReplicas не формирует список отставших, и needsShrinkIsr сообщает, что сжатие не требуется: пока предыдущее изменение не завершено, новый список отставших не формируется.
Обновление состояния реплики при fetch-запросе
При каждом fetch-запросе под read-блокировкой leaderIsrUpdateLock лидер обновляет состояние реплики: смещение, время последнего fetch и эпоху брокера. Обновлённая эпоха затем участвует в допуске в ISR — isReplicaIsrEligible сравнивает сохранённую эпоху с кэшированной. Если фолловер пришлёт fetch со старой эпохой брокера, сохранённая эпоха не совпадёт с кэшированной, и реплика в ISR не попадёт. Помимо эпохи, реплика допускается только если она не fenced и не в controlled shutdown.
Как формируется и отправляется AlterPartition
Лидер не может менять ISR сам — он лишь предлагает новое состояние контроллеру через AlterPartitionзапрос лидера к контроллеру с предложением нового состава ISR. Два способа формирования различаются тем, что лидер предполагает о результате.
При расширении (prepareIsrExpand) новый синхронный реплика добавляется к текущему ISR, и лидер сразу исходит из того, что тот войдёт в ISR до подтверждения контроллера, — так HW успеет отразить новый ISR даже при задержке подтверждения. Сохраняется ожидающее состояние PendingExpandIsr с новым LeaderAndIsr. При сжатии (prepareIsrShrink) из ISR исключаются вышедшие из синхронизации реплики, но успех обновления не предполагается: если бы лидер продвинул HW в расчёте на сжатие, а AlterPartition не прошёл, HW оказался бы завышен. Поэтому «максимальный ISR» для PendingShrinkIsr — это текущий ISR. Оба метода формируют LeaderAndIsr и присваивают partitionState = updatedState, то есть состояние становится ожидающим подтверждения контроллера.
При успешном подтверждении состояние заменяется на CommittedPartitionState с ISR и leaderRecoveryStateпризнаком состояния восстановления лидера для партиции из ответа, эпоха партиции обновляется, и вызывается markIsrExpand или markIsrShrink в зависимости от типа ожидающего состояния.
Отправка AlterPartition вынесена из-под блокировки лидерства, потому что логика завершения может увеличить HW и тем самым завершить отложенные операции. При успешном ответе после обновления состояния вызывается maybeIncrementLeaderHW; если HW вырос, вызывается tryCompleteDelayedRequests, который завершает все отложенные операции. Этот метод документирован как вызываемый без удержания leaderIsrUpdateLock.
Ошибки обрабатываются по-разному. Для UNKNOWN_TOPIC_OR_PARTITION, UNKNOWN_TOPIC_ID, FENCED_LEADER_EPOCH, INVALID_UPDATE_VERSION, INVALID_REQUEST и NEW_LEADER_ELECTED повтор не выполняется: состояние остаётся ожидающим, а лидер ждёт новых метаданных. Для OPERATION_NOT_ATTEMPTED и INELIGIBLE_REPLICA повтора тоже нет, но перед этим partitionState сбрасывается к последнему закоммиченному состоянию — контроллер явно сообщает, что запрос не был применён. Для всех прочих ошибок запрос отправляется повторно.
При успешном ответе handleAlterPartitionUpdate сначала проверяет, что эпоха лидера в ответе совпадает с текущей, а эпоха партиции не меньше текущей; иначе ответ игнорируется и помечается неудача. Устаревшая эпоха лидера относится к другому лидеру, меньшая эпоха партиции — к более старой версии состояния, которую у нас уже перекрыли.
Как лидер вычисляет high watermark
high watermark. maybeIncrementLeaderHW начинается с проверки isUnderMinIsr: если партиция ниже минимума ISR, HW не увеличивается. Иначе кандидат инициализируется как LEO лидера, а затем для каждой удалённой реплики его LEO понижает кандидата, если он меньше текущего кандидата и при этом либо brokerId входит в maximalIsr, либо истинно shouldWaitForReplicaToJoinIsr — условие, которое выделяет реплику, догоняющую лидера и имеющую право войти в ISR: она считается догнавшей относительно LEO лидера и одновременно проходит isReplicaIsrEligible. Учёт таких реплик при выборе нового HW и есть механизм, не дающий догоняющей реплике застрять навсегда: её LEO участвует в вычислении HW ещё до того, как она официально войдёт в ISR.
maximalIsr включает реплики, добавленные через AlterPartition, но ещё не подтверждённые контроллером. Учёт их делает набор более ограничительным, а значит безопасным: даже если контроллер откатит ISR, HW не окажется завышенным. В конце вызывается maybeIncrementHighWatermark, который обновляет HW только если новое значение больше старого (или лежит на более новом сегменте), иначе HW не меняется.
Учёт догоняющих реплик нужен, чтобы избежать ловушки: HW вычисляется как минимум LEO по репликам, которые либо в ISR, либо считаются догнавшими. Если бы HW продвигался только по LEO лидера (когда ISR состоит из одного лидера), LEO отстающей реплики оставался бы меньше растущего HW, условие followerEndOffset >= leaderLog.highWatermark никогда бы не выполнялось, и реплика не попала бы в ISR никогда. Комментарий прямо описывает этот эффект: без ожидания фолловера его LEO может вечно отставать от HW и он никогда не будет добавлен в ISR.
Почему проверяется именно недостижение минимума, а не непустота ISR: HW может двигаться только когда размер ISR не меньше min ISR. isUnderMinIsr истинно, когда размер ISR меньше effectiveMinIsr и локальная реплика — лидер, а effectiveMinIsr — это минимум из настроенного minInSyncReplicas и числа реплик (удалённые плюс лидер). Так при потере реплик ниже порога min.insync.replicas лидер не подтверждает новые сообщения через HW, даже если ISR ещё не пуст.
Ограничения HW в журнале и усечение
updateHighWatermark ограничивает новое значение с двух сторон: снизу — logStartOffset (если переданное смещение меньше, берётся logStartOffset), сверху — конец лога (если переданное смещение не меньше конца, берётся конец). При немонотонном понижении (новое значение меньше текущего) пишется предупреждение, но значение всё равно присваивается. После этого уведомляются менеджер состояния продюсеров и слушатель смещений, и вызывается обновление первого нестабильного смещения.
Усечение лога при возврате отставшей реплики устроено не по HW, а по эпохе лидераLeader Epoch — 32-битном монотонно возрастающем номере непрерывного периода лидерства для партиции. Лидер проставляет эпоху на всех сообщениях, а каждая реплика хранит вектор [LeaderEpoch => StartOffset], отмечающий смены лидеров в истории её лога; этот вектор заменяет HW при усечении.
Фолловер посылает лидеру LeaderEpochRequest для своей текущей эпохи. Лидер возвращает либо конец лога, если следующей эпохи нет, либо первый офсет следующей эпохи, — то есть ответ содержит офсет, на котором заканчивается запрошенная эпоха. Фолловер усекает лог до этого офсета, удаляя только те сообщения, которых нет в логе лидера.
Два сценария показывают, почему усечение по HW опасно. В первом фолловер успел получить сообщение, но не успел обновить HW; он перезапускается, усекает лог до HW, затем лидер падает, и фолловер становится лидером — сообщение потеряно навсегда. Суть проблемы в том, что фолловеру нужен дополнительный раунд RPC для обновления HW, и этот разрыв позволяет быстрой смене лидера привести к потере закоммиченного сообщения. Во втором после сбоя питания обе машины расходятся из-за асинхронного сброса на диск: если лидером окажется машина с наименьшим логом, данные теряются, а реплики расходятся с разной историей. При поддержке сжатых наборов сообщений расхождение в худшем случае вообще останавливает репликацию — офсет сжатого набора в одной реплике может указывать в середину сжатого набора в другой. Решение по эпохам устраняет оба: в первом случае ответ возвращает офсет 2 вместо HW = 0, и сообщение не усекается; во втором после перезапуска лидер принимает новое сообщение с новой эпохой, а вернувшийся фолловер получает первый офсет этой эпохи, понимает, что его сообщение осиротело, и усекает его.
Пограничные случаи и что видит клиент
При фехтованииfenced — пометке брокера как недоступного для лидерства контроллер сначала удаляет брокера из ISR и выбирает новых лидеров для затронутых партиций, а затем фиксирует запись о фехтовании. При контролируемом завершении работы, если брокер ещё не в этом состоянии, сначала фиксируется запись о нём, после чего брокер также удаляется из ISR и выбираются новые лидеры. Для партиций без лидера, где unclean leader election разрешён для топика, контроллер выбирает лидера нечистым способом, добавляя записи о выборе; признак необходимости немедленно повторить операцию — достижение числа записей порога maxElectionsPerImbalance, который ограничивает, сколько выборов лидера выполняется за один проход, чтобы не перегружать контроллер и дать ему возможность вернуться к остальным разделам на следующей итерации.
ISR становится пустым, когда все брокеры-реплики партиции отфильтрованы как неактивные — то есть помечены как fenced или находящиеся в controlled shutdown.
Когда отставший брокер возвращается и становится фолловером, он отправляет запрос эпох, чтобы определить точку усечения. Лидер сравнивает переданную эпоху с локальной: если локальная больше — возвращается FENCED_LEADER_EPOCH, если меньше — UNKNOWN_LEADER_EPOCH, иначе ошибки нет. На стороне фолловера способ усечения выбирается по наличию эпохи: для партиций с известной эпохой выполняется усечение по границе эпохи, для партиций без эпохи — усечение до HW, где целевым офсетом служит текущий fetch-офсет. При FENCED_LEADER_EPOCH партиция помечается как ошибочная; при прочих ошибках она также помечается для повторной попытки.
Пока идёт усечение и последующая догонка, HW фолловера остаётся на прежнем уровне. Потребитель видит только сообщения до HW: updateHighWatermark ограничивает HW сверху концом лога, а чтение HW возвращает текущее значение. То есть клиент не увидит ни данных, которые лидер ещё не подтвердил, ни данных, которые фолловер вот-вот усечёт.
Отдельно про согласованность ISR: лидер не может менять его сам. Контроллер отклоняет AlterPartition, если предложенный ISR не содержит текущего лидера, и если запрос пришёл не от текущего лидера. Дополнительно контроллер помечает негодные реплики: незарегистрированный брокер, брокер в controlled shutdown, fenced брокер, брокер с несовпадающей эпохой; при непустом списке таких реплик возвращается INELIGIBLE_REPLICA. При этой ошибке или OPERATION_NOT_ATTEMPTED лидер откатывает состояние к последнему закоммиченному, поскольку контроллер явно сообщает, что запрос не был применён.
Что из этого следует на практике
- HW — это не «сколько записал лидер», а минимум по репликам, которые либо в ISR, либо догнали и имеют право в него войти. Поэтому подтверждение продюсеру приходит только после того, как данные забрал достаточный набор реплик.
- Реплика попадает в ISR не по числовому лагу, а по двум границам: HW лидера и начало текущей эпохи лидера. Числовой порог
replicaLagTimeMaxMsиспользуется только при исключении из ISR, а не при включении. - Без учёта догоняющих реплик в расчёте HW отставшая реплика не смогла бы догнать растущий HW и никогда не вернулась бы в ISR. Поэтому HW намеренно придерживается на LEO догоняющей реплики.
- Пока изменение ISR не подтверждено контроллером, лидер при расширении считает нового реплику частью ISR (безопасно, набор строже), а при сжатии — нет (иначе HW мог бы уйти за пределы реально закоммиченного).
- Усечение лога делается по эпохе лидера, а не по HW. Это устраняет окно, в котором фолловер, не успевший обновить HW, усекает закоммиченное сообщение и теряет его при быстрой смене лидера.
- Потребитель видит только данные до HW, поэтому ни неподтверждённые сообщения лидера, ни усекаемые сообщения вернувшегося брокера клиенту не видны.
Где смотреть в коде
- Partition.scala: canAddReplicaToIsr
- Partition.scala: handleAlterPartitionError
- Partition.scala: maybeIncrementLeaderHW
- Partition.scala: handleAlterPartitionUpdate
- Partition.scala: needsShrinkIsr
- UnifiedLog.java: fetchHighWatermarkMetadata
- UnifiedLog.java: updateHighWatermark
- Partition.scala: prepareIsrShrink