Назад к блогу

Как ISR и high watermark защищают Kafka от потери сообщений при отказе брокера

Как ISR и high watermark защищают Kafka от потери сообщений при отказе брокера

Когда брокер Kafka падает или возвращается с отставшим логом, судьба подтверждённых сообщений зависит от того, как лидер определяет границу коммита и состав синхронных реплик. В статье разбирается механика high watermark и ISR: какие проверки проходят реплики, чтобы попасть в набор синхронных, и как принимается решение об их исключении. Отдельное внимание — асинхронной природе этих решений на лидере и контроллере и цене ошибок: от потери уже подтверждённых данных до бесконечного отставания реплики.

В Kafka сообщение считается закоммиченным не тогда, когда его записал лидер, а тогда, когда его успели забрать достаточно реплик. Границей этой «достаточности» служит high watermark, а множеством реплик, которые обязаны успеть, — ISR. Вопрос, который разбираем: как именно лидер решает, кого держать в ISR, как вычисляется HW и что происходит, когда брокер выпадает из строя или возвращается с отставшим логом. Нетривиальность в том, что все эти решения принимаются асинхронно и распределённо, а цена ошибки — либо потеря подтверждённых данных, либо вечная остановка реплики, которая никак не может догнать растущий HW.

Кто попадает в ISR и как из него выпадают

ISR. Лидер проверяет три условия по порядку, с коротким замыканием: по реплике не идёт незавершённая операция изменения ISR (состояние inflight), её идентификатора ещё нет в текущем ISR, и она проходит проверку isReplicaIsrEligible. Порядок такой, потому что более дешёвые проверки состояния отсекают реплику раньше, чем выполняется обращение к метаданным брокера.

isReplicaIsrEligible определяет, имеет ли реплика право войти в ISR. Сначала она достаёт реплику из состояния лидера; если её там нет (например, топик уже удалён), реплика в ISR не допускается. Затем проверяются три вещи: реплика не fenced, не находится в controlled shutdown, и её сохранённая эпоха брокера совпадает с кэшированной в метаданных. Совпадение эпох проверяется так:

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. Два способа формирования различаются тем, что лидер предполагает о результате.

При расширении (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, а по эпохе лидера. Лидер проставляет эпоху на всех сообщениях, а каждая реплика хранит вектор [LeaderEpoch => StartOffset], отмечающий смены лидеров в истории её лога; этот вектор заменяет HW при усечении.

Фолловер посылает лидеру LeaderEpochRequest для своей текущей эпохи. Лидер возвращает либо конец лога, если следующей эпохи нет, либо первый офсет следующей эпохи, — то есть ответ содержит офсет, на котором заканчивается запрошенная эпоха. Фолловер усекает лог до этого офсета, удаляя только те сообщения, которых нет в логе лидера.

Два сценария показывают, почему усечение по HW опасно. В первом фолловер успел получить сообщение, но не успел обновить HW; он перезапускается, усекает лог до HW, затем лидер падает, и фолловер становится лидером — сообщение потеряно навсегда. Суть проблемы в том, что фолловеру нужен дополнительный раунд RPC для обновления HW, и этот разрыв позволяет быстрой смене лидера привести к потере закоммиченного сообщения. Во втором после сбоя питания обе машины расходятся из-за асинхронного сброса на диск: если лидером окажется машина с наименьшим логом, данные теряются, а реплики расходятся с разной историей. При поддержке сжатых наборов сообщений расхождение в худшем случае вообще останавливает репликацию — офсет сжатого набора в одной реплике может указывать в середину сжатого набора в другой. Решение по эпохам устраняет оба: в первом случае ответ возвращает офсет 2 вместо HW = 0, и сообщение не усекается; во втором после перезапуска лидер принимает новое сообщение с новой эпохой, а вернувшийся фолловер получает первый офсет этой эпохи, понимает, что его сообщение осиротело, и усекает его.

Пограничные случаи и что видит клиент

При фехтовании контроллер сначала удаляет брокера из 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, поэтому ни неподтверждённые сообщения лидера, ни усекаемые сообщения вернувшегося брокера клиенту не видны.

Где смотреть в коде

Источники