Назад к блогу

Cluster Mirroring в Kafka: репликация между кластерами силами брокеров

Cluster Mirroring в Kafka: репликация между кластерами силами брокеров

В Kafka появился встроенный механизм репликации между кластерами — Cluster Mirroring (KIP-1279), который переносит данные силами самих брокеров, без внешних Connect-воркеров. Функциональность пока в фазе Early Access и управляется feature gate, но уже заявлена как замена MirrorMaker 2 с его проблемами перевода смещений и обработки нечистых выборов лидера.

В Kafka появился механизм Cluster Mirroring (KIP-1279) — репликация между кластерами, встроенная в сами брокеры, без внешних Connect-воркеров. Зеркало описывается как именованный однонаправленный канал от удалённого исходного кластера к локальному целевому, а данные перекладываются между кластерами по fetch-протоколу с сохранением смещений и сжатия.

Функциональность управляется feature gate mirror.version и на момент написания находится в фазе Early Access: по умолчанию она выключена, а все mirroring API отклоняются с ошибкой UNSUPPORTED_VERSION.

Что именно изменилось

Зеркало кластера задаётся именем, bootstrap-серверами источника и учётными данными безопасности. Bootstrap-серверы — обязательная конфигурация bootstrap.servers, список пар хост/порт для начального соединения с исходным кластером. Поддерживаются стандартные механизмы аутентификации Kafka, включая TLS/SSL и SASL; каждое зеркало может иметь собственные настройки безопасности. Все свойства и учётные данные хранятся как ресурс конфигурации CLUSTER_MIRROR в логе метаданных целевого кластера и используются только для аутентифицированных соединений с источником. Ресурс конфигурации — это именованная запись настроек в логе метаданных, по образцу конфигураций брокеров и топиков: свойства зеркала задаются при его создании или изменении и хранятся в метаданных кластера, а не в файлах на диске брокера, поэтому доступны всему кластеру и переживают перезапуски. Тот же ресурс задействован в административных запросах на чтение и изменение конфигураций.

Управлять зеркалами можно через Admin API и через CLI-инструмент kafka-cluster-mirrors.sh с флагами --create, --delete, --describe, --failed, --json, --list, --mirror, --mirror-config, --pause, --recover, --resume, --start, --stop, --topics:

bin/kafka-cluster-mirrors.sh --bootstrap-server :9094 --create --mirror a-to-b --mirror-config /tmp/a-to-b.properties

Внутреннее состояние зеркал хранится в служебном топике __mirror_state: туда координатор пишет записи для отслеживания состояния зеркала и точек синхронизации между брокерами.

Состояния партиции и эпохи

Состоянием зеркальной партиции владеет конечный автомат, который сохраняет переходы состояний в __mirror_state и управляет жизненным циклом репликации: усечением лога, fencing эпох, failover.

Состояния различаются тем, что происходит с репликацией:

  • MIRRORING — создаётся поток mirror fetcher для выборки из исходного кластера;
  • PAUSED — ничего не делается;
  • LOG_ALIGNMENT — инициируется процесс усечения;
  • FAILED — наступает, в частности, при несовпадении ID исходного кластера во время периодического обновления метаданных, а также при нечистом выборе лидера на источнике, если поддержка такой ситуации выключена;
  • ULE_RECOVERY — наступает при нечистом выборе лидера на источнике, если поддержка такой ситуации включена; ожидание идёт до сходимости всех назначенных реплик, а не только ISR, к усечённому смещению, и потеря данных на источнике принимается.

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

Предыстория: почему не MirrorMaker 2

В KIP-1279 перечислены пять проблем MirrorMaker 2:

  • Lossy Offset Translation. Перевод смещений «inherently lossy»: MM2 не может гарантировать возврат той же самой записи, а при отсутствии точного перевода гарантирует лишь, что запись по переведённому смещению всегда раньше фактической, что ведёт к дублирующей обработке при failover.
  • External Offset Complexity. Платформы вроде Flink и Spark, source-коннекторы Kafka Connect и транзакционные приложения хранят смещения вне __consumer_offsets и при failover вынуждены обращаться к внутреннему топику offset-sync MM2, что добавляет операционную сложность и точки отказа.
  • Unclean Leader Election. MM2 не обрабатывает нечистые выборы лидера на источнике корректно: потеря данных на источнике может не отразиться на назначении, и состояние кластеров расходится.
  • Operational Burden. MM2 работает как отдельные Connect-воркеры вне брокеров, требуя отдельного развёртывания, мониторинга и управления жизненным циклом, включая независимые от брокеров обновления.
  • Compression Inefficiency. Если записи источника сжаты, MM2 распаковывает и сжимает их заново при записи в назначение, что снижает пропускную способность и увеличивает задержку.

Как это работает

Репликация идёт напрямую по fetch-протоколу: целевой брокер сам забирает данные из исходного кластера consumer Fetch-запросами, а не follower Fetch. Причина в том, что лидер целевого кластера не зарегистрирован как follower в исходном кластере: follower Fetch заставил бы исходный брокер обновлять статус реплики, о которой он не знает. Consumer Fetch таких побочных эффектов не несёт, при этом сохраняются те же проверки согласованности логов, что и при внутрикластерной репликации.

В состоянии MIRRORING сжатые record batches копируются побайтово и добавляются в локальный лог без повторного сжатия, с сохранением исходного типа сжатия (gzip, snappy, lz4, zstd или без него). В сводном сравнении MM2 отмечен как «Decompress and recompress», а Cluster Mirroring — «Byte-for-byte passthrough».

Смещения на источнике и назначении совпадают, включая пропуски, оставленные компактификацией топика. Поэтому потребительским группам не нужна трансляция смещений: зафиксированное смещение на источнике равно зафиксированному смещению на назначении. Идентичность смещений опирается на общий topic ID: источник и назначение разделяют один и тот же topic ID, что позволяет fetch-запросам проходить валидацию на исходном брокере и отличать тот же логический топик от совпадения имён. При наличии MirrorInfo с topic ID контроллер проверяет, что ID действителен, не занят другим именем топика в целевом кластере и что все реплики в назначении партиций активны. Если topic ID равен ZERO_UUID (источник на pre-2.8 Kafka без поддержки topic IDs), контроллер генерирует случайный UUID, система переходит на поиск топика по имени и пропускает более строгую валидацию реплик.

На каждом брокере работают три компонента:

  • Менеджер метаданных зеркал периодически синхронизирует метаданные из исходных кластеров: состояние топиков, конфигурации, смещения групп, ACL. Хранит кэш зеркальных партиций, заполняемый воспроизведением __mirror_state-партиций, которые лидирует этот брокер.
  • Координатор зеркал владеет конечным автоматом партиции, сохраняет переходы в __mirror_state и управляет жизненным циклом репликации.
  • Менеджер фетчеров создаёт и ведёт потоки, которые реплицируют данные с исходных брокеров.

Менеджер фетчеров организует потоки по трёхмерному ключу: идентификатор фетчера, endpoint исходного брокера, имя зеркала. Партиции разных зеркал получают отдельные потоки для изоляции аутентификации, партиции одного зеркала распределяются по нескольким потокам для балансировки. Смена лидера в исходной партиции вызывает переназначение или пересоздание потока на нового исходного брокера.

Кэш зеркальных партиций живёт в координаторе зеркал. При получении лидерства над __mirror_state лог воспроизводится для перестроения кэша, при потере лидерства кэш очищается. Когда брокер становится лидером зеркальной партиции данных, он определяет текущее состояние партиции: если этот же брокер координирует нужную партицию __mirror_state, состояние читается из локального in-memory кэша, иначе состояние запрашивается у брокера-координатора и ответ ожидается асинхронно. Дальше по состоянию выбирается действие: при MIRRORING создаётся поток фетчера, при PAUSED не делается ничего, при LOG_ALIGNMENT запускается усечение.

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

Периодическая синхронизация метаданных идёт с каждым брокером: он получает описания топиков из источника, обновляет кэши лидеров, обнаруживает удалённые топики и восстанавливает пропущенные партиции. Создание топиков, масштабирование партиций, конфигурации топиков, смещения потребительских и share-групп, ACL, повышение эпохи лидера, обнаружение топиков по include-шаблону и применение exclude-шаблона выполняются только на координаторе. Для потребительских групп зафиксированные смещения периодически забираются из источника и применяются на назначении; для share-групп из источника извлекается текущий SPSO (Share-Partition Start Offset) и применяется на назначении, что инициализирует состояние группы и в координаторе групп, и в share-координаторе. Интервал по умолчанию — mirror.metadata.refresh.interval.ms, 60000 мс. Во время обновления брокер проверяет, что ID исходного кластера не изменился; при несовпадении синхронизация для этого зеркала останавливается и переводится в терминальное состояние FAILED.

Двухфазное усечение

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

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

Включение функциональности

Функциональность управляется feature gate mirror.version. При значении 0 (по умолчанию) все mirroring API отклоняются с UNSUPPORTED_VERSION, при 1 — включаются. Feature gate требует в качестве минимума новую версию метаданных.

В фазе Early Access для включения все узлы кластера — контроллеры и брокеры — должны явно задать unstable.api.versions.enable=true и unstable.feature.versions.enable=true во всех конфигурационных файлах. После запуска кластера с минимальной версией метаданных операторы включают mirror.version динамически:

bin/kafka-features.sh --bootstrap-server :9092 upgrade --feature mirror.version=1

Что это меняет на практике

Failover запускается командой остановки зеркалирования — StopMirrorTopics. После её завершения зеркальные топики переходят из read-only в writable, продюсеры переподключаются с новыми producer ID и sequence numbers, а консьюмеры продолжают с последних синхронизированных смещений под тем же group ID. Failback — создание нового зеркала на старом исходном кластере, где назначением становится новый источник, и запуск обратного зеркалирования; если исходный кластер поддерживает Cluster Mirroring, зеркалируется только дельта.

Миграция кластера идёт по шагам: поднять новый KRaft-кластер целевой версии, создать зеркало со старого кластера на новый, запустить зеркалирование (назначение само обнаруживает топики, создаёт совпадающие партиции с идентичными topic IDs, синхронизирует конфигурации и начинает репликацию данных), остановить продюсеров на старом кластере, следить за replication lag до нуля, остановить зеркалирование (топики на новом кластере становятся writable), перенаправить клиентов и вывести старый кластер из эксплуатации.

RPO зависит от replication lag на момент сбоя: зеркалирование асинхронное, поэтому записи, произведённые на источнике, но ещё не реплицированные на назначение, при незапланированном failover будут потеряны.

Ограничения и открытые вопросы

Авторы признают, что exactly-once semantics между кластерами не поддерживается: единственная гарантия — отсутствие зависших транзакций, блокирующих READ_COMMITTED-потребителей после failover, атомарная синхронизация зафиксированных записей не гарантируется. Mirror fetcher использует READ_UNCOMMITTED, поэтому данные незафиксированных транзакций видны на приёмнике до failover; приложениям со строгими транзакционными гарантиями рекомендуется дедупликация или сверка после failover.

Интеграция с tiered storage не поддерживается: во время LOG_ALIGNMENT локальный лог усекается до LME, что несовместимо с tiered storage на приёмнике, а при включённом tiered storage для mirror-топика его партиции переходят в FAILED. Diskless Topics на момент написания ещё обсуждаются (KIP-1500 и другие под-KIP), поддержка появится в будущих KIP.

Active-active топология изначально не поддерживается; двунаправленное зеркалирование возможно только для разных топиков, а при обнаружении пересечения затронутые партиции переводятся в failed state.

Throttling реализован с двух сторон: на приёмнике — broker-level mirror.replication.throttled.rate и topic-level mirror.replication.throttled.replicas, на источнике — через клиентские квоты по детерминированному client identifier каждого mirror fetcher, кодирующему ID брокера, номер потока фетчера и имя зеркала.

Миграция с MirrorMaker 2 плавным переходом невозможна: системы используют разные внутренние структуры топиков и по-разному хранят смещения потребителей. Переключение состоит из остановки репликации MM2, удаления зеркальных топиков на целевом кластере вместе с внутренними топиками MM2 и запуска Cluster Mirroring с нуля.

По масштабируемости при 100k партиций на брокер in-memory кэш, хранящий состояние партиции, эпоху и метаданные отказов, остаётся небольшим по heap — десятки МБ. Основное давление создают периодическое обновление метаданных (по умолчанию каждые 60 с) с батчевыми запросами к исходному кластеру за описаниями топиков, конфигурациями и смещениями потребительских групп, ответы на которые на таком масштабе становятся большими, и обработка крупных дельт при переназначении лидерства, когда MMM строит промежуточные коллекции пропорционально размеру дельты и рассылает координаторам запросы на чтение состояния зеркальных партиций. Первое можно смягчить увеличением интервала обновления под размер развёртывания; во втором случае размер запроса растёт пропорционально числу партиций, одновременно меняющих лидерство. Оба случая могли бы выиграть от последующего KIP с chunked-запросами или лимитами размера RPC.

Источники

Похожее