Назад к блогу

Apache Kafka: нативная репликация кластеров без внешних инструментов

Apache Kafka: нативная репликация кластеров без внешних инструментов

Новый механизм Cluster Mirroring (KIP-1279), встроенный прямо в брокеры Kafka, позволяет реплицировать данные между кластерами без внешних инструментов вроде MirrorMaker 2 и Kafka Connect. Репликация работает на уровне fetch-протокола, сохраняя байтовую идентичность батчей, точные смещения и поддержку unclean leader election — всё это упрощает архитектуру и снижает накладные расходы.

Зачем нужна нативная репликация кластеров

Три отдельных сервиса, которые нужно синхронизировать между собой, — именно столько инфраструктуры обычно стоит за кросс-кластерной репликацией Kafka. Cluster Mirroring (KIP-1279) убирает эти сервисы: механизм встроен прямо в брокеры Kafka. Destination-брокер сам забирает зафиксированные записи у source-кластера через тот же fetch-протокол, которым пользуются обычные follower-реплики, и добавляет их в локальный лог партиции байт в байт.

Что это меняет по сравнению с MirrorMaker 2 (MM2) — стандартным инструментом с времён Kafka 2.4, который работает как набор воркеров Kafka Connect (consume из источника, produce в приёмник):

  • Нет внешней инфраструктуры. Потоки mirror fetcher работают внутри процесса брокера, а не в отдельном кластере Connect.
  • Байтовая передача. Сжатые батчи переносятся как raw bytes — без распаковки и повторного сжатия. Батч gzip, snappy, lz4 или zstd приходит в приёмник в исходном виде; CPU-затраты на полный цикл decompress/recompress исчезают, а выбор компрессии продюсера сохраняется.
  • Точное сохранение смещений. Лог приёмника держит те же offset'ы, что и источник, включая пропуски от компакции топиков. Consumer group переключается без трансляции: зафиксированное смещение на источнике — это то же смещение на приёмнике. В MM2 это были потери при трансляции через топик.
  • Переключение одной командой. Флаг --stop проводит партицию по детерминированной последовательности: fetcher'ы удаляются, последняя mirror-эпоха сохраняется в LME, leader epoch инкрементируется, незавершённые транзакции абортируются, а новый control record истекает состояние всех продюсеров. Партиция становится доступной на запись — без внешней координации и запросов offset'ов.
  • Поддержка unclean leader election. Если на источнике прошли ULE, приёмник входит в состояние восстановления и ждёт сходимости всех назначенных реплик, а не только ISR, к усечённому смещению, прежде чем возобновить репликацию. Согласованность логов сохраняется даже когда источник выбрал лидера с неполным логом.

Итоговая разница между инструментами:

АспектMirrorMaker 2Cluster mirroring
РазвёртываниеВнешние воркеры ConnectВнутри брокера
КомпрессияРаспаковка и повторное сжатиеБайтовый passthrough
СмещенияПотери при трансляции через топикИдентичны между кластерами
Переключение консьюмеровЗапрос к топику синхронизации offset'овНапрямую, без трансляции
Unclean electionsНе обрабатываютсяПолная сходимость логов
Совместимость с источникамиKafka 2.0+Kafka 2.1+
МониторингСпецифичные для Connect инструментыСтандартные JMX-метрики брокера

Репликация — не единственное, что берёт на себя брокер. Он также отвечает за обнаружение метаданных, синхронизацию конфигураций, смещения consumer group и распространение ACL. Управление полосой пропускания работает с обеих сторон: приёмник применяет настраиваемый лимит скорости репликации, а на источнике mirror-трафик выглядит как обычные запросы консьюмеров — существующие квоты клиентов применяются без изменений.

Как Cluster Mirroring превосходит MirrorMaker 2

Начнём с развертывания. MM2 — это набор воркеров Kafka Connect: чтобы репликация заработала, нужен отдельный кластер Connect с коннектором MirrorSourceConnector. Cluster Mirroring обходится без него — потоки mirror fetcher живут внутри процесса брокера, поэтому внешний контур разворачивать не нужно.

С обработкой данных разница не менее ощутима. MM2 распаковывает батч, а затем сжимает его снова — CPU тратится на каждый батч, а исходный выбор компрессии продюсера теряется. Cluster Mirroring передаёт сжатые батчи как raw bytes: gzip, snappy, lz4 или zstd приходит на приёмник в том же виде, в котором был создан.

Управление смещениями — третий параметр, и здесь MM2 заметно проигрывает. Он ведёт отдельный топик синхронизации смещений, а перенос из source в destination идёт через перевод (lossy translation). Cluster Mirroring хранит те же смещения, что и источник, включая пропуски от компакции топиков, поэтому смещение коммита на источнике равно смещению на приёмнике — перевод не нужен вообще. При смене кластера клиент просто подключается к новому bootstrap-серверу и продолжает с того же места.

Отказоустойчивость — четвёртый параметр. Если на источнике случается unclean leader election (ULE), MM2 ничего с этим не делает. Cluster Mirroring с включённым mirror.unclean.leader.election.enable переводит партицию в состояние ULE_RECOVERY: она ждёт, пока все назначенные реплики, а не только участники ISR, догонят усечённый offset, и лишь затем возобновляет репликацию. При значении false (по умолчанию) партиция уходит в FAILED и требует ручного вмешательства.

Остальные отличия сведены в таблицу:

ПараметрMirrorMaker 2Cluster Mirroring
РазвертываниеВнешние воркеры ConnectВстроено в брокер
КомпрессияРаспаковка и повторное сжатиеБайтовый проброс
СмещенияLossy translation через топикИдентичны между кластерами
Failover потребителейЗапрос к топику синхронизации смещенийНапрямую, без перевода
Нечистые выборыНе обрабатываютсяПолная сходимость логов
Совместимость с источникомKafka 2.0+Kafka 2.1+
МониторингИнструменты ConnectСтандартные JMX-метрики брокера

Помимо самих данных, брокер-приёмник выполняет обнаружение метаданных, синхронизацию конфигураций топиков, перенос смещений групп и распространение ACL. Ограничение пропускной способности настраивается с обеих сторон: приёмник применяет лимит скорости репликации, а трафик mirror fetch на источнике выглядит как обычные consumer-запросы, так что штатные квоты клиентов работают без изменений.

Внутренние компоненты Cluster Mirroring

Репликация запускается одной командой, но за кулисами работают три независимых механизма. Разберём каждый — и то, как они связаны.

MirrorMetadataManager (MMM) — это оркестратор, и он живёт на каждом брокере. MMM реализует интерфейс MetadataPublisher, то есть подписывается на изменения в KRaft-журнале метаданных. Когда контроллер записывает MirrorTopicStateChangeRecord, MMM на лидере затронутой партиции запускает нужный переход: create, start, stop, pause, resume, recover или delete. Плюс MMM держит Admin-клиент к исходному кластеру и по умолчанию раз в 60 секунд обновляет его метаданные: находит новые топики по include/exclude-шаблонам, синхронизирует конфигурации, забирает смещения consumer-групп и проверяет, что ID исходного кластера не изменился. Последняя проверка не декоративная: если кто-то перенастроит зеркало на другой кластер, без неё данные поехали бы молча и с порчей.

ClusterMirrorCoordinator (CMC) отвечает за хранение состояния. Он устроен по тому же паттерну, что group coordinator и transaction coordinator: управляет шардами внутреннего компактированного топика __mirror_state (по умолчанию cleanup.policy=compact, 50 партиций, фактор репликации 3). Состояние каждой mirror-партиции лежит в этом топике как пара ключ-значение, а конкурентный доступ разграничен optimistic concurrency control через leader epoch и state epoch.

MirrorFetcherThread (MFT) делает основную работу. Он наследуется от AbstractFetcherThread — того же базового класса, что используется для репликации внутри кластера, — и тянет записи из источника, добавляя их в локальные логи. У каждого потока свой NetworkClient с отдельными учётными данными под конкретное зеркало, поэтому SASL/SSL-контексты разных зеркал не пересекаются. Менеджер фетчеров индексирует потоки по трём измерениям — ID фетчера, endpoint исходного брокера и имя зеркала, — что даёт тонкую балансировку нагрузки и быструю реакцию на смену лидера на источнике.

Связь между ними такая: MMM ловит изменение состояния и запускает переход, CMC фиксирует результат в __mirror_state, а MFT выполняет фактическую передачу данных. Пример: администратор вызывает --pause. MMM видит запись в метаданных и инициирует переход в PAUSING, CMC сохраняет новое состояние партиции, а MFT сворачивает поток фетчера — партиция остаётся доступной только для чтения. При --resume новый фетчер поднимается сразу в MIRRORING, без повторного выравнивания логов.

Жизненный цикл зеркальной партиции: состояния и переходы

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

LOG_ALIGNMENT: выравнивание локального лога

Первое состояние после старта — выравнивание локального лога с источником. Здесь всё зависит от наличия записи LME (Last Mirror Epoch), которая хранит последнее зеркальное смещение и эпоху из предыдущей сессии. Если LME нет — например, зеркалирование запускается впервые или исходный кластер не поддерживает нужные механизмы — брокер обрезает лог до нуля и копирует данные заново. Если LME найдена, зеркалится только дельта: брокер обрезает записи лишь на смещениях вплоть до последнего зеркального смещения, разрешая любое расхождение.

После выравнивания идёт EPOCH_FENCING: брокер отправляет контроллеру запрос BumpLeaderEpochs, увеличивая локальную лидерскую эпоху на 10, с порогом повторного увеличения равным 3. Гарантия критична — без неё потребитель на стороне назначения мог бы инициализироваться с зафиксированной эпохой из источника, превышающей локальную, и отвергнуть лидера как устаревшего.

MIRRORING: рабочий режим

Поток-выборщик начинает тянуть записи из источника, а сжатые батчи добавляются в локальный лог без распаковки и повторного сжатия. High watermark продвигается по мере репликации на локальные фолловеры, а смещения consumer-групп периодически синхронизируются из источника и обрезаются до допустимого диапазона на стороне назначения.

ULE_RECOVERY: реакция на нечистые выборы

После запуска зеркалирования расхождение логов может возникнуть только из-за unclean leader election в исходном кластере — ведь зеркалятся исключительно зафиксированные записи. Если поддержка ULE включена и лидер назначения замечает расхождение при выборке, партиция переходит из MIRRORING в ULE_RECOVERY. Выборщик удаляется, и система ждёт сходимости всех реплик, а не только членов ISR, к обрезанному концу лога. Как только каждая реплика догнала — партиция возвращается в MIRRORING, создаётся новый выборщик.

Настройка mirror.unclean.leader.election.enable даёт явный выбор. При true партиция уходит в ULE_RECOVERY и ждёт сходимости всех назначенных реплик, принимая потерю данных от ULE источника. При false (значение по умолчанию) партиция переходит в FAILED, полностью останавливая зеркалирование.

Пауза, остановка и повтор

Оператор может приостановить зеркалирование без потери прогресса: потоки-выборщики сносятся, но партиция остаётся доступной только для чтения. Возобновление возвращает непосредственно в MIRRORING со свежими выборщиками, минуя повторное выравнивание.

Остановка — это путь к переключению при отказе. Брокер удаляет выборщики, записывает LME для будущего обратного переключения, увеличивает локальную эпоху, добавляет ABORT-маркеры для незавершённых транзакций и пишет управляющую запись MIRROR_PID_RESET, которая истекает всё состояние производителей. После этого партиция становится доступной для записи. Перезапуск остановленного зеркала проходит снова через LOG_ALIGNMENT, поскольку за время работы на запись в топике могли накопиться локальные данные.

FAILED: восстановление после ошибок

Ошибки запускают автоматический повтор с экспоненциальной задержкой и джиттером — до настраиваемого максимума попыток. Непортящиеся ошибки требуют ручного вмешательства: например, смена ID исходного кластера или удаление топика-источника.

Механизм сохранения согласованности эпох

Потребитель на принимающем кластере приходит со своим зафиксированным смещением и эпохой лидера. Если локальная эпоха окажется меньше исходной, брокер отклонит лидера как устаревшего, и клиент зависнет. Отсюда жёсткое правило: эпоха лидера на приёмнике всегда больше или равна эпохе лидера на источнике (DLE >= SLE). Разрыв поддерживается тремя способами: реактивный подъём, когда эпоха очередного полученного батча подходит близко к локальной; упреждающий, когда зазор сжимается меньше 3; и периодический, во время синхронизации метаданных источника. Каждый подъём увеличивает эпоху на 10, а запрос BumpLeaderEpochs уходит контроллеру.

Расхождение логов решает двухфазный протокол усечения. Начальная фаза выравнивает истории эпох: запись LME хранит последнее зеркальное смещение и эпоху прошлой сессии, и во время LOG_ALIGNMENT брокер обрезает записи только на смещениях вплоть до этого последнего зеркального смещения. Если LME нет, зеркалирование идёт с нуля. Устойчивая фаза выравнивает смещения внутри оставшихся эпох уже в ходе обычного фетча: поток-выборщик определяет версию fetch-запроса источника на лету. При Fetch v12+ информация о расходящихся эпохах в ответе позволяет усечь лог инлайново, ровно до точки расхождения. Если версия ниже, усечение при фетче отключается, и поток отдельно запрашивает OffsetsForLeaderEpoch, чтобы найти ту же точку. Так схема работает и со старыми кластерами-источниками, переживая обновления и откаты версий.

Раз во время зеркалирования принимаются только зафиксированные записи, любое усечение после LOG_ALIGNMENT может прийти лишь от неаккуратных выборов лидера (ULE) на источнике. Флаг mirror.unclean.leader.election.enable даёт явный выбор. Со значением true партиция уходит в ULE_RECOVERY и ждёт, пока сойдутся все назначенные реплики, а не только члены ISR, — и вы сознательно принимаете потерю данных от ULE на источнике. Со значением false (по умолчанию) партиция падает в FAILED, и зеркалирование останавливается целиком. Выбор простой: либо продолжаете с возможной потерей, либо сохраняете целостность ценой остановки.

Детали протокола и настройки — в KIP-1279.

Практические сценарии: аварийное восстановление и миграция

Если основной кластер недоступен, счёт идёт на минуты, а не на количество выполненных шагов из runbook. Пока кластер A обслуживает продакшн-трафик, зеркало a-to-b непрерывно копирует топики в кластер B. Для этого достаточно создать и запустить зеркало:

kafka-cluster-mirrors.sh --bootstrap-server B:9092 \
  --create --mirror a-to-b --mirror-config mirror.properties

kafka-cluster-mirrors.sh --bootstrap-server B:9092 \
  --start --mirror a-to-b --topics ".*"

Состояние репликации в любой момент показывает флаг --describe. Когда кластер A падает, вы останавливаете зеркало — и топики на B становятся доступными для записи:

kafka-cluster-mirrors.sh --bootstrap-server B:9092 \
  --stop --mirror a-to-b

Продюсеры и консьюмеры переключают bootstrap-серверы на B. Поскольку смещения идентичны, а смещения консьюмер-групп синхронизированы, приложения продолжают с того места, где остановились: ни пересчёта смещений, ни повторной обработки. Переключение занимает одну команду, потому что --stop последовательно снимает фетчеры, фиксирует LME, поднимает эпоху лидера, добавляет ABORT-маркеры для незавершённых транзакций и записывает MIRROR_PID_RESET, обнуляя состояние продюсеров.

После восстановления A вы поднимаете обратное зеркало:

kafka-cluster-mirrors.sh --bootstrap-server A:9092 \
    --create --mirror b-to-a --mirror-config reverse.properties

kafka-cluster-mirrors.sh --bootstrap-server A:9092 \
    --start --mirror b-to-a --topics ".*"

Система узнаёт, что A уже был источником для этих топиков, находит сохранённый при failover LME и обрезает лог инкрементально вместо полной перерепликации. Когда репликация догоняет, вы выполняете --stop --mirror b-to-a и возвращаете клиентов на A. Помните про границу: зеркалирование асинхронное, так что записи, не успевшие уехать из A в B, при незапланированном отказе теряются. Для большинства сценариев DR это осознанный компромисс — синхронный режим добавил бы задержку основному кластеру. Нулевой RPO обещают в KIP-1360.

Миграция выигрывает от совместимости источников вплоть до Kafka 2.1. Вместо последовательного обновления 2.x → 3.x → 4.x с миграцией ZooKeeper в KRaft на месте вы поднимаете новый KRaft-кластер, создаёте и запускаете зеркало со старого:

kafka-cluster-mirrors.sh --bootstrap-server new:9092 \
    --create --mirror old-to-new --mirror-config old-cluster.properties

kafka-cluster-mirrors.sh --bootstrap-server new:9092 \
    --start --mirror old-to-new --topics ".*"

Приёмник сам обнаруживает топики, создаёт партиции с теми же topic ID, подтягивает конфигурации и начинает репликацию. Останавливаете продюсеров на старом кластере, ждёте, пока лаг дойдёт до нуля, и завершаете миграцию:

kafka-cluster-mirrors.sh --bootstrap-server new:9092 \
    --stop --mirror old-to-new

После этого клиенты переезжают на новый кластер, а старый можно выводить из эксплуатации. Ни промежуточных обновлений версий, ни конвертации метаданных на месте, ни скриптов миграции ZooKeeper-to-KRaft: старый кластер работает без изменений, пока вы сами не решите его погасить.

Ограничения и будущее развитие

Три вещи, которые стоит унести из материала. Первое: репликация между кластерами перестала быть внешним процессом — зеркалирование встроено в брокер через KIP-1279, поэтому сжатые батчи идут на приёмник байт в байт, смещения совпадают, а консьюмер-группы переключаются без таблиц трансляции. Второе: миграция больше не требует последовательного апгрейда через каждую промежуточную версию — источником может быть кластер вплоть до Kafka 2.1, а новый KRaft-кластер поднимается параллельно и получает данные зеркалом. Третье, самое важное для планирования: зеркалирование асинхронное, и RPO напрямую зависит от задержки репликации в момент отказа. Записи, которые продюсер успел отправить в источник, но которые не дошли до приёмника, при незапланированном переключении теряются. Синхронный режим только заявлен в KIP-1360, а синхронная репликация между кластерами добавила бы задержку основному кластеру — для большинства сценариев аварийного восстановления на это не идут.

Практический совет: не полагайтесь на --describe вручную в момент аварии — сделайте задержку репликации метрикой с алертом. Регулярно сверяйте отставание приёмника от источника и держите порог, после которого вы либо останавливаете продюсеров на источнике, либо признаёте потерю приемлемой. При плановом переключении — например, при миграции — просто дождитесь нулевой задержки перед --stop, и тогда RPO равен нулю без всякой синхронности. Источник: Data liberation: Apache Kafka's native cluster mirroring.

Похожее