В 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 в логе метаданныхжурнал изменений состояния кластера, который ведут контроллеры KRaft целевого кластера и используются только для аутентифицированных соединений с источником. Ресурс конфигурации — это именованная запись настроек в логе метаданных, по образцу конфигураций брокеров и топиков: свойства зеркала задаются при его создании или изменении и хранятся в метаданных кластера, а не в файлах на диске брокера, поэтому доступны всему кластеру и переживают перезапуски. Тот же ресурс задействован в административных запросах на чтение и изменение конфигураций.
Управлять зеркалами можно через 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внешний инструмент репликации Kafka на базе Connect-воркеров:
- 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режим работы Kafka, в котором метаданные кластера хранятся в логе метаданных, который ведут контроллеры, а не в ZooKeeper-кластер целевой версии, создать зеркало со старого кластера на новый, запустить зеркалирование (назначение само обнаруживает топики, создаёт совпадающие партиции с идентичными 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.