Зачем нужна нативная репликация кластеров
Три отдельных сервиса, которые нужно синхронизировать между собой, — именно столько инфраструктуры обычно стоит за кросс-кластерной репликацией 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 2 | Cluster 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 2 | Cluster 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.