Изучаем Apache Kafka с нуля. Урок 30. MirrorMaker: репликация между кластерами

Изучаем Apache Kafka с нуля. Урок 30. MirrorMaker: репликация между кластерами

 

В уроке 29 мы разобрали kafka-consumer-perf-test.sh — как измерить скорость чтения, интерпретировать вывод и подбирать параметры консьюмера. Получили полную картину: есть цифры продюсера, есть цифры консьюмера, узкое место теперь видно.

Следующий логичный вопрос: а что делать, если кластеров несколько? DR-окружение, географически распределённые датацентры, изоляция окружений prod и staging — всё это требует переноса данных между кластерами Kafka. Инструмент для этого называется MirrorMaker.

Теория репликации между кластерами и практика настройки production-сценариев подробно разбирается в курсе «Администрирование кластера Kafka» — там же настройка мониторинга репликационного лага и стратегии failover.

 

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

Внутренняя репликация Kafka — то, что мы уже разбирали в уроках про partition и replication factor — работает внутри одного кластера. Она защищает от падения отдельных брокеров, но не от потери всего датацентра. Если кластер целиком недоступен, внутренняя репликация не поможет.

Репликация между кластерами решает другие задачи. Во-первых, disaster recovery: есть основной кластер и резервный в другом датацентре, данные синхронизируются, и при сбое трафик переключается на резервный. Во-вторых, агрегация: несколько региональных кластеров сливают данные в один центральный. В-третьих, изоляция: prod-топики реплицируются в отдельный кластер для аналитики или QA без прямого доступа к production.

Для всех этих сценариев и создан MirrorMaker — утилита, которая читает сообщения из одного кластера и пишет их в другой.

 

MirrorMaker 1 и почему его больше нет в Kafka 4.x

Исторически первая версия MirrorMaker (MM1) запускалась скриптом kafka-mirror-maker.sh. По сути это был обычный консьюмер + продюсер в одном процессе: читаем из source-кластера, пишем в target. Просто, понятно, работало.

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

В Kafka 2.7 MM1 официально объявили устаревшим. В Kafka 4.0 скрипт kafka-mirror-maker.sh был удалён полностью. Если вы работаете с Kafka 4.x — его в директории bin/ нет, и это нормально. Вместо него используется MirrorMaker 2, который запускается через connect-mirror-maker.sh.

 

Что такое MirrorMaker 2

MirrorMaker 2 (MM2) — это набор коннекторов для Kafka Connect, а не отдельный инструмент. Он построен на том же фреймворке, что и все коннекторы к внешним системам. Репликация данных здесь — это просто коннектор, который читает из одного кластера Kafka и пишет в другой.

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

MM2 решает все проблемы первой версии. Он синхронизирует offset’ы consumer-групп между кластерами — при переключении на резервный кластер консьюмеры продолжают с той же позиции. Он автоматически создаёт реплицируемые топики в target-кластере с нужными параметрами. Он масштабируется горизонтально: несколько воркеров Kafka Connect делят нагрузку по партициям.

Архитектура MM2 включает три компонента. MirrorSourceConnector реплицирует данные и метаданные топиков из source в target. MirrorCheckpointConnector синхронизирует offset’ы consumer-групп. MirrorHeartbeatConnector отправляет heartbeat-сообщения для мониторинга задержки репликации.


flowchart LR
    subgraph source["Source Cluster"]
        ST["Топики\n(orders, payments)"]
        CG["Consumer\nGroups"]
    end

    subgraph mm2["MirrorMaker 2\n(Kafka Connect Workers)"]
        SC["MirrorSource\nConnector"]
        CC["MirrorCheckpoint\nConnector"]
        HC["MirrorHeartbeat\nConnector"]
    end

    subgraph target["Target Cluster"]
        TT["source.orders\nsource.payments"]
        TO["Синхронизированные\noffset'ы групп"]
        TH["source.heartbeats"]
    end

    ST --> SC --> TT
    CG --> CC --> TO
    SC --> HC --> TH

Файл конфигурации mm2.properties для MirrorMaker 2

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

# Объявляем псевдонимы кластеров
clusters = source, target

# Адреса брокеров
source.bootstrap.servers = source-broker:9092
target.bootstrap.servers = target-broker:9092

# Включаем репликацию из source в target
source->target.enabled = true

# Какие топики реплицировать (регулярное выражение)
source->target.topics = .*

# Синхронизация offset'ов consumer-групп
source->target.sync.group.offsets.enabled = true
source->target.sync.group.offsets.interval.seconds = 60

# Какие consumer-группы включить в синхронизацию
source->target.groups = .*

# Фактор репликации для внутренних топиков MM2
# (поставьте 3 в production)
replication.factor = 1

# Количество задач для каждого коннектора
tasks.max = 4

Несколько параметров требуют отдельного внимания.

source->target.topics принимает регулярное выражение. Значение .* реплицирует все топики. Если нужны только конкретные — укажите паттерн: orders.*|payments. Внутренние топики MM2 в репликацию не попадают по умолчанию.

replication.factor здесь отвечает за внутренние служебные топики MirrorMaker, а не за реплицируемые данные. Фактор репликации целевых топиков наследуется из исходного кластера, если не указать source->target.replication.factor явно.

Топики в target-кластере по умолчанию получают префикс с именем source-кластера: топик orders в source становится source.orders в target. Это сделано специально, чтобы избежать коллизий при двусторонней репликации. Если префикс мешает, можно его переопределить через replication.policy.class.

 

Запуск MirrorMaker 2 кластера с connect-mirror-maker.sh

После того как файл конфигурации готов, запуск выглядит просто:

# Запуск MirrorMaker 2 в standalone-режиме
connect-mirror-maker.sh /etc/kafka/mm2.properties

По умолчанию MM2 запускается в режиме distributed Kafka Connect — он сам создаёт внутренние топики для хранения конфигурации и статуса коннекторов. Это требует, чтобы оба кластера были доступны на момент старта.

Для производственного запуска лучше передать дополнительные параметры JVM и указать файл с настройками воркера отдельно:

# Запуск с настройками heap и отдельным файлом воркера
export KAFKA_HEAP_OPTS="-Xms512m -Xmx2g"

connect-mirror-maker.sh \
  /etc/kafka/mm2.properties

# Запуск в фоне с логированием
connect-mirror-maker.sh /etc/kafka/mm2.properties \
  >> /var/log/kafka/mirror-maker.log 2>&1 &

MM2 создаёт в target-кластере несколько служебных топиков сразу после старта: mm2-offsets.source.internal, mm2-status.source.internal, mm2-configs.source.internal, а также source.heartbeats и source.checkpoints.internal. Их наличие — первый признак того, что MM2 запустился корректно.

 

Проверка работы репликации

После запуска нужно убедиться, что данные действительно переезжают. Первый шаг — проверить, что топики появились в target-кластере:

# Список топиков в target-кластере
kafka-topics.sh \
  --bootstrap-server target-broker:9092 \
  --list

# Ожидаемый вывод:
# source.orders
# source.payments
# source.heartbeats
# mm2-offsets.source.internal
# mm2-status.source.internal

Затем проверяем, что heartbeat-топик наполняется. MM2 пишет туда сообщения с заданным интервалом, и задержка между ними показывает репликационный лаг:

# Читаем несколько heartbeat-сообщений
kafka-console-consumer.sh \
  --bootstrap-server target-broker:9092 \
  --topic source.heartbeats \
  --from-beginning \
  --max-messages 5 \
  --property print.timestamp=true

Offset’ы consumer-групп проверяются через kafka-consumer-groups.sh на target-кластере. Если sync.group.offsets включён, группы из source должны появиться и там — это значит, что failover сработает без потери позиции.

Альтернативы

MM2 через connect-mirror-maker.sh — стандартный open-source путь. Но есть и другие варианты. Confluent Replicator — коммерческий коннектор от Confluent с расширенными возможностями мониторинга и управления схемами. uReplicator от Uber решает проблему балансировки при большом числе топиков. Brooklin от LinkedIn поддерживает репликацию не только между кластерами Kafka, но и между разными типами очередей сообщений. Для простых сценариев с небольшим числом топиков иногда достаточно настроить продюсер-консьюмер цепочку вручную, особенно если нужна трансформация данных на лету.

Схема Confluent repliocator для иллюстрации концепции MirrorMaker 2 Apache Kafka

 

Apache Kafka: администрирование кластера

Код курса
KAFKA
Ближайшая дата курса
5 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

В качестве заключения

 

MirrorMaker 2 — это набор Kafka Connect коннекторов, и в этом уроке мы запустили его в standalone-режиме через connect-mirror-maker.sh. В уроке 31 разберём connect-standalone.sh подробно: как запускать произвольные коннекторы к внешним системам, как устроен worker, как читается конфигурация коннектора. Это откроет полную картину экосистемы Kafka Connect.

Референсные ссылки

Все уроки курса

Тема URL
1 Установка Kafka с Zookeeper https://test.bigdataschool.ru/blog/news/lesson1-kafka-zookeeper-install/
2 Установка Kafka в режиме KRaft https://test.bigdataschool.ru/blog/news/lesson2-kafka-kraft-install/
3 Docker KRaft: однонодовый кластер https://test.bigdataschool.ru/blog/news/lesson3-kafka-docker-single/
4 Docker KRaft: 3-нодовый кластер https://test.bigdataschool.ru/blog/news/lesson4-kafka-docker-cluster/
5 Утилиты bin/: переменные окружения и основы https://test.bigdataschool.ru/blog/news/lesson5-kafka-bin-intro/
6 kafka-topics.sh: управление топиками https://test.bigdataschool.ru/blog/news/lesson6-kafka-topics/
7 kafka-console-producer.sh https://test.bigdataschool.ru/blog/news/lesson7-kafka-console-consumer/
8 kafka-console-consumer.sh https://test.bigdataschool.ru/blog/news/lesson8-kafka-console-consumer/
9 kafka-server-start.sh / kafka-server-stop.sh https://test.bigdataschool.ru/blog/news/lesson9-kafka-server-start-stop/
10 kafka-storage.sh https://test.bigdataschool.ru/blog/news/lesson10-kafka-storage/
11 Шпаргалка по кластеру Kafka https://test.bigdataschool.ru/blog/news/lesson11-kafka-cluster-cheatsheet/
12 kafka-configs.sh (broker) https://test.bigdataschool.ru/blog/news/lesson12-kafka-configs-broker/
13 kafka-metadata-shell.sh https://test.bigdataschool.ru/blog/news/lesson13-kafka-metadata-shell/
14 kafka-features.sh https://test.bigdataschool.ru/blog/news/lesson14-kafka-features/
15 kafka-configs.sh (topics & clients) https://test.bigdataschool.ru/blog/news/lesson15-kafka-configs/
16 kafka-log-dirs.sh https://test.bigdataschool.ru/blog/news/lesson16-kafka-log-dirs/
17 kafka-dump-log.sh https://test.bigdataschool.ru/blog/news/lesson17-kafka-dump-log/
18 kafka-delete-records.sh https://test.bigdataschool.ru/blog/news/lesson18-kafka-delete-records/
19 kafka-consumer-groups.sh https://test.bigdataschool.ru/blog/news/lesson19-kafka-consumer-groups/
20 kafka-streams-application-reset.sh https://test.bigdataschool.ru/blog/news/lesson20-kafka-streams-reset/
21 kafka-leader-election.sh https://test.bigdataschool.ru/blog/news/lesson21-kafka-leader-election/
22 kafka-reassign-partitions.sh https://test.bigdataschool.ru/blog/news/lesson22-kafka-reassign-partitions/
23 kafka-replica-verification.sh https://test.bigdataschool.ru/blog/news/lesson23-kafka-replica-verification/
24 kafka-acls.sh https://test.bigdataschool.ru/blog/news/lesson24-kafka-acls/
25 kafka-broker-api-versions.sh https://test.bigdataschool.ru/blog/news/lesson25-kafka-broker-api-versions/
26 kafka-get-offsets.sh https://test.bigdataschool.ru/blog/news/lesson26-kafka-get-offsets/
27 kafka-verifiable-producer/consumer.sh https://test.bigdataschool.ru/blog/news/lesson27-kafka-verifiable/
28 kafka-producer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson28-kafka-producer-perf/
29 kafka-consumer-perf-test.sh https://test.bigdataschool.ru/blog/news/lesson29-kafka-consumer-perf/
30 kafka-mirror-maker.sh https://test.bigdataschool.ru/blog/news/lesson30-kafka-mirror-maker/
31 connect-standalone.sh https://test.bigdataschool.ru/blog/news/lesson31-kafka-connect-standalone/
32 connect-distributed.sh https://test.bigdataschool.ru/blog/news/lesson32-kafka-connect-distributed/
33 kcat. Альтернативный CLI https://test.bigdataschool.ru/blog/news/lesson33-kcat/