MirrorMaker 2 и геораспределённый Kafka: репликация между кластерами, DR, active-active

Один кластер Kafka живёт в одном дата-центре — а бизнесу нужны катастрофоустойчивость и присутствие в нескольких регионах. Разбираем два подхода: растянутый (stretch) кластер против репликации отдельных кластеров через MirrorMaker 2; что MM2 на самом деле переносит (топики, offset'ы, ACL, конфиги) и как работает трансляция offset при failover; топологии active-passive/active-active/hub-and-spoke; RPO/RTO, failover и failback; цена латентности и трафика между регионами. С живым стендом из двух KRaft-кластеров и MM2.

Пока Kafka живёт в одном дата-центре, всё просто: репликация между брокерами, ISR, acks=all — и данные переживут падение узла. Но что если падает весь регион? Или сервис обязан работать в нескольких регионах сразу? Один кластер этого не решает — начинается территория геораспределённого Kafka: либо растянуть кластер на несколько площадок, либо держать отдельные кластеры и реплицировать между ними. У каждого пути своя цена, и MirrorMaker 2 (MM2) — основной инструмент для второго.

Девятая, финальная статья серии «Kafka: глубокое погружение». Опирается на модель репликации и ISR из статьи про надёжность и эксплуатацию из статьи про операции. MM2 построен поверх Kafka Connect — механику коннекторов оттуда переиспользуем. Числа ниже — с живого стенда из двух KRaft-кластеров (us-east, 3 брокера — тот же кластер-фундамент, что и в предыдущих статьях серии, и us-west, single-node) и реального MM2. Конфигурация MM2 и стенд из двух кластеров — в репозитории digital-cookbook, каталог kafka/mm2, а оркестрирующий сценарий (репликация, трансляция offset, failover, active-active) — скрипт ops/mm2-demo.sh.

MirrorMaker 2 и гео-репликация: два кластера, трансляция offset и failover без потерь

В статье

Stretch-кластер против репликации кластеров

Есть два принципиально разных способа сделать Kafka устойчивым к падению целого дата-центра или региона.

Stretch (растянутый) кластер — один логический кластер, брокеры и контроллеры которого физически размещены в разных зонах или дата-центрах. Механика та же, что и в обычном кластере, разобранная в статье про ISR и репликацию: broker.rack задаёт зону для каждого брокера, rack-aware назначение реплик старается не класть все копии партиции в одну зону, acks=all по-прежнему ждёт подтверждения от всего текущего ISR — просто теперь часть этого ISR физически находится за десятки или сотни километров. Это даёт синхронную, сильную консистентность: подтверждённая запись гарантированно пережила отказ одной площадки. Цена — задержка каждого acks=all-коммита равна RTT до самой дальней реплики ISR, а при недоступности одной из зон кластер должен либо ждать (durability), либо рисковать (unclean leader election, разобранный в той же статье). На практике это ограничивает stretch-кластер зонами доступности в пределах одного региона или метро-расстояниями с низким и стабильным RTT — растягивать синхронный кворум на разные континенты означает добавлять десятки-сотни миллисекунд в каждую запись.

Репликация отдельных кластеров — противоположный компромисс: два (или больше) полностью независимых кластера, каждый со своим ISR и своим контроллерным кворумом, между которыми асинхронно копируются данные специальным инструментом. Локальная запись в каждом кластере не ждёт ничего за пределами своего региона — латентность производителя не меняется от того, что где-то на другом конце света есть вторая копия. Плата за это — eventual-консистентность между кластерами: копия на другой стороне всегда немного отстаёт, и при внезапной потере исходного кластера часть данных, ещё не успевших реплицироваться, теряется. MirrorMaker 2 — стандартный инструмент Apache Kafka для этой второй модели.

Выбор между ними — не про «что лучше», а про то, какой RTT можно себе позволить в пути записи. Синхронный stretch — для метро-расстояний, где задержка приемлема и консистентность важнее локальной скорости. MM2 — для настоящих гео-расстояний (другой континент, разные облачные регионы) и для сценариев disaster recovery, где сама идея ждать межконтинентальный ISR на каждой записи неприемлема.

Как устроен MirrorMaker 2

MM2 — не отдельный протокол и не скрипт поверх консольного продюсера/консьюмера, а набор коннекторов Kafka Connect, развёрнутых как distributed workers. Это значит, что вся механика Connect — офсеты воркера, статус задач, масштабирование через tasks.max — применяется к MM2 напрямую.

Три коннектора решают три разные задачи:

  • MirrorSourceConnector — читает записи топика на source-кластере и пишет их на target, сохраняя партиционирование и порядок внутри партиции.
  • MirrorCheckpointConnector — периодически фиксирует соответствие между offset consumer-группы на source и её позицией на target (checkpoint), а при включённом sync.group.offsets ещё и напрямую пишет транслированный committed offset в __consumer_offsets на target.
  • MirrorHeartbeatConnector — регулярно пишет heartbeat-записи, по которым можно судить о том, что канал репликации вообще жив, а не просто перестал получать новые данные из-за того, что source-топик молчит.

Реплицированный топик получает имя с префиксом кластера-источника — us-east.orders, а не просто orders. Это не косметика, а единственный реальный механизм защиты от циклов репликации в active-active-конфигурации: если оба кластера реплицируют друг в друга под одинаковым явным списком топиков (не regex-маской), топик с префиксом целевого кластера уже не совпадает с фильтром исходного имени и не попадает под повторную репликацию назад. У MM2 есть и встроенная эвристика (DefaultReplicationPolicy распознаёт префикс и по умолчанию исключает уже помеченные топики из списка кандидатов на репликацию), но полагаться только на неё рискованно — там, где нужна предсказуемость, явный список топиков в конфиге надёжнее скрытого поведения.

Что MM2 переносит: данные топика, consumer offset’ы (через checkpoint), ACL (если sync.topic.acls.enabled и оба кластера используют авторизацию) и конфиги топиков (при sync.topic.configs). Чего не переносит автоматически: переключение самих клиентов на другой кластер — MM2 подготавливает данные и offset’ы к failover, но кто именно из producer’ов и consumer’ов когда переключит bootstrap.servers, решает приложение или оператор, а не MM2.

Живой пример: репликация топика orders

Стенд — два независимых KRaft-кластера: us-east (3-брокерный кластер-фундамент серии, использованный и в предыдущих статьях) и us-west (новый single-node KRaft-кластер, поднятый специально для этого стенда — одного брокера достаточно, потому что предмет демонстрации здесь механика MM2, а не повторная отказоустойчивость target-кластера, уже разобранная в статье про ISR и репликацию). MM2 запущен в dedicated-режиме (connect-mirror-maker.sh), подключён сразу к обеим docker-сетям.

clusters = us-east, us-west

us-east.bootstrap.servers = kafka1:9092,kafka2:9092,kafka3:9092
us-west.bootstrap.servers = west1:9092

# --- Направление репликации ---
us-east->us-west.enabled = true
us-west->us-east.enabled = false

# --- Явный фильтр: только наш демо-топик и наша демо-группа ---
# (не ".*" по умолчанию — на us-east уже живут топики предыдущих стендов серии)
us-east->us-west.topics = orders
us-east->us-west.groups = orders-cg.*

# --- Трансляция consumer offset'ов ---
emit.checkpoints.enabled = true
sync.group.offsets.enabled = true
sync.group.offsets.interval.seconds = 5

# ACL-синхронизация выключена: оба кластера PLAINTEXT без авторизации
sync.topic.acls.enabled = false

На us-east создан топик orders (3 партиции, RF=3, min.insync.replicas=2), в него отправлено 200 сообщений с acks=all — все 200 отправлены и подтверждены (sent=200 acked=200). MM2 с refresh.topics.interval.seconds=5 подхватил новый топик и реплицировал его на us-west под именем us-east.orders: 200 из 200 записей, число партиций сохранено 1:1, порядок внутри каждой партиции сохранён.

Отдельно стоит отметить: MM2 на apache/kafka 4.3.1 в этом стенде поднялся с первой попытки — все три коннектора (source/checkpoint/heartbeat) стартовали одновременно, без единого ERROR или Exception в логе за весь прогон.

Трансляция consumer offset и failover

Самая содержательная часть MM2 — не копирование данных (это тривиально), а перенос позиции consumer-группы так, чтобы после переключения на другой кластер она не начинала читать с нуля и не теряла место. Проблема в том, что offset — это номер записи в конкретной партиции конкретного топика конкретного кластера: у топика orders на us-east и топика us-east.orders на us-west нумерация абсолютно независимая, потому что это два разных лога, наполняемых в разное время и, в общем случае, не синхронно байт-в-байт.

На стенде группа orders-cg прочитала 120 из 200 записей orders на us-east (CURRENT-OFFSET=120, LOG-END-OFFSET=200, LAG=80) и закоммитила эту позицию. MirrorCheckpointConnector с sync.group.offsets.enabled=true (интервал 5 секунд) сам, без явного вызова RemoteClusterUtils.translateOffsets, перенёс committed offset той же группы на us-west — но не как 120, а как 102.

Разрыв в 18 не ошибка и не потеря точности: MM2 фиксирует соответствие offset’ов source↔target периодически (через внутренний топик offset-syncs), и трансляция намеренно консервативна вниз — гарантия «не пропустить непрочитанное» важнее гарантии «точное число». Если бы трансляция округляла вверх, консьюмер после failover рисковал бы пропустить записи, которые формально ещё не читал; округление вниз означает, что часть уже прочитанных записей будет прочитана повторно (at-least-once), но ни одна непрочитанная не потеряется.

Дальше — реальная симуляция failover, а не имитация: consumer той же группы orders-cg запущен на us-west против топика us-east.orders без --from-beginning. Он стартовал с транслированного offset 102 и дочитал оставшиеся 98 записей (102…199) до CURRENT-OFFSET=200, LAG=0.

sum_group_col() {
  # sum_group_col BROKER GROUP TOPIC COL — суммирует колонку COL
  # (4=CURRENT-OFFSET, 5=LOG-END-OFFSET, 6=LAG) kafka-consumer-groups.sh
  # --describe по ВСЕМ партициям топика — распределение записей по
  # партициям не гарантировано, суммировать нужно честно, не только
  # партицию 0.
  docker exec "$1" /opt/kafka/bin/kafka-consumer-groups.sh \
    --bootstrap-server localhost:9092 --describe --group "$2" \
    | awk -v t="$3" -v c="$4" '$2==t{sum+=$c; n++} END{print sum+0, n+0}'
}

# после failover-консьюмера: суммарный LAG по ВСЕМ партициям должен == 0
read -r lag_sum lag_n <<< "$(sum_group_col kafka-cookbook-west-1 orders-cg "us-east.orders" 6)"
[ "$lag_n" -gt 0 ] && [ "$lag_sum" -eq 0 ]  # [assert] OK

Итог проверен суммой по всем партициям, а не предположением «всё лежит в партиции 0»: 18 записей (102…119) прочитаны повторно, 80 новых (120…199) прочитаны впервые — 100% покрытие исходных 200 записей, ноль потерь, ноль пропусков. Проверка воспроизведена и на «грязном» топике с неравномерным распределением по партициям (700 записей, 200/500/0): суммарный LAG после failover снова сошёлся в ноль, то есть сам механизм не зависит от того, как записи разложились по партициям.

А вот величину разрыва обобщать не стоит: 18 — число конкретного прогона, а не константа механизма. Разрыв определяется тем, когда именно MirrorCheckpointConnector успел зафиксировать соответствие offset’ов в offset-syncs относительно момента коммита группы, поэтому при другом тайминге он окажется другим. Устойчиво здесь ровно одно свойство, и его достаточно: транслированный offset никогда не превышает исходный committed.

sequenceDiagram participant App as Consumer (orders-cg) participant East as us-east: orders participant MM2 as MirrorCheckpointConnector participant West as us-west: us-east.orders App->>East: читает и коммитит offset=120 Note over East: CURRENT-OFFSET=120, LAG=80 MM2->>East: наблюдает committed offset MM2->>West: sync.group.offsets: пишет транслированный offset=102 Note over West: разрыв 18 — консервативная трансляция вниз Note over App,East: аварийная недоступность us-east App->>West: FAILOVER: стартует БЕЗ --from-beginning West-->>App: offset=102 (не 0, не 120) App->>West: дочитывает 102..199 (98 записей) Note over West: LAG=0 по всем партициям, 18 повторов, 0 потерь

sequenceDiagram
    participant App as Consumer (orders-cg)
    participant East as us-east: orders
    participant MM2 as MirrorCheckpointConnector
    participant West as us-west: us-east.orders

    App->>East: читает и коммитит offset=120
    Note over East: CURRENT-OFFSET=120, LAG=80
    MM2->>East: наблюдает committed offset
    MM2->>West: sync.group.offsets: пишет транслированный offset=102
    Note over West: разрыв 18 — консервативная трансляция вниз

    Note over App,East: аварийная недоступность us-east
    App->>West: FAILOVER: стартует БЕЗ --from-beginning
    West-->>App: offset=102 (не 0, не 120)
    App->>West: дочитывает 102..199 (98 записей)
    Note over West: LAG=0 по всем партициям, 18 повторов, 0 потерь
Failover consumer-группы: транслированный offset вместо чтения с начала

Лаг самой репликации данных — величина измеримая, но, как и все throughput-числа в этой серии, host-зависимая. В трёх прогонах 300 дополнительных записей на us-east (acks=all) занимали на продюсирование 1711, 1945 и 2076 мс, а us-west догонял паритет (полные 500 записей суммарно) через 1.9, 4.2 и 4.5 секунды после завершения продюсирования — разброс догона более чем вдвое на одной и той же машине показывает, насколько бессмысленно цитировать здесь одно «характерное» число. Абсолютные миллисекунды не значат ничего сами по себе (один docker-хост, без реальной межрегиональной сети); значим сам факт конечного, измеримого лага — прямое следствие того, что MM2 репликация асинхронная, а не синхронная.

Топологии: active-passive, active-active, hub-and-spoke

Active-passive (DR) — самый распространённый сценарий: продакшн-нагрузка целиком идёт в primary-кластер, MM2 непрерывно льёт копию в standby. В обычном режиме standby не обслуживает клиентов вообще — он существует ради одного момента, аварии primary, когда клиенты переключаются на него, подхватывая транслированные offset’ы вместо чтения с начала. Именно этот сценарий и разобран выше на живом примере.

Active-active — оба кластера принимают запись от собственных локальных клиентов, и MM2 реплицирует в обе стороны одновременно. На стенде это проверено буквально: локальный топик orders создан прямо на us-west (не как реплика, а как собственные данные региона), MM2 переключён на конфиг с обоими направлениями репликации включёнными.

clusters = us-east, us-west

us-east->us-west.enabled = true
us-west->us-east.enabled = true

us-east->us-west.topics = orders
us-east->us-west.groups = orders-cg.*

us-west->us-east.topics = orders
us-west->us-east.groups = .^

50 сообщений, отправленных локально на us-west, реплицировались на us-east как us-west.orders — 50 из 50. На обеих сторонах при этом не появилось ни us-west.us-east.orders, ни us-east.us-west.orders: цикл не запустился. Честно: в этом конкретном стенде реальным барьером был точный список топиков (topics = orders, не regex .*) с обеих сторон одновременно — имя orders не совпадает с уже префиксованным us-east.orders, поэтому повторная репликация назад не находит, что реплицировать. Встроенная эвристика MM2 по распознаванию префикса — дополнительная подстраховка, а не единственный работающий здесь механизм.

Активно-активная топология требует готовности приложения к тому, чего в single-кластерной Kafka не бывает: между двумя независимо принимающими запись кластерами нет глобального порядка и нет общей exactly-once-семантики. Если один и тот же логический объект обновляется параллельно с двух сторон, конфликт разрешает приложение (last-write-wins по временной метке, CRDT-подобные структуры, партиционирование по ключу так, чтобы конкретный объект всегда писался только с одной стороны) — MM2 такой конфликт не видит и не решает.

Hub-and-spoke (fan-in/fan-out) — третий паттерн, не проверенный на этом стенде живьём, но логически прямое расширение той же механики: несколько региональных кластеров реплицируют данные в центральный (fan-in, типичный случай — аналитика или единый data lake, собирающий события со всех регионов), либо, наоборот, центральный кластер раздаёт конфигурацию или референсные данные во все региональные (fan-out). Технически это тот же набор из трёх коннекторов MM2, просто настроенных на несколько пар кластеров одновременно, а не на одну.

DR: RPO/RTO, failover и failback

Для disaster recovery поверх MM2 важны две метрики. RPO (Recovery Point Objective, сколько данных допустимо потерять) определяется лагом асинхронной репликации в момент отказа — тем самым лагом, измеренным выше (~1.9 секунды в характерном прогоне на локальном стенде, на реальной межрегиональной сети заведомо больше). Всё, что producer успел записать на primary, но MM2 ещё не успел скопировать на target, при внезапной полной потере primary теряется безвозвратно. RTO (Recovery Time Objective, за сколько времени сервис снова доступен) определяется процедурой переключения: обнаружение отказа, переключение клиентов на target, проверка, что consumer-группы подхватили транслированные offset’ы.

Failover — переключение продюсеров и консьюмеров на target-кластер. Именно это отработано на живом стенде выше: consumer стартует на us-west без --from-beginning, берёт транслированный offset вместо чтения с начала или потери позиции, и LAG сходится к нулю ровно так, как если бы переключения кластера не было.

Failback — обратное движение, вернуться на восстановленный primary после того, как авария устранена, — и на практике это обычно сложнее failover. Пока приложение работало на target, туда успели прийти новые записи (в active-passive DR — от временно переключённых producer’ов, если такое допускалось; в active-active — по определению топологии). Чтобы вернуться на primary, не потеряв эти записи, MM2 должен доливать их обратно тем же механизмом (target становится временным source), и здесь легко задвоить данные, если процедура не учитывает уже присутствующие на primary записи до аварии. Правильная процедура failback — это по сути ещё один цикл репликации и трансляции offset, только в обратном направлении, выполняемый осознанно и с проверкой, а не автоматически.

Важно сказать честно: на этом стенде полная потеря source-региона (реальный docker kill всего кластера us-east) и процедура failback не воспроизводились живьём. Причина не техническая сложность, а то, что us-east — общий кластер-фундамент всей серии, и он должен был оставаться исправным для остальных статей. Механика, описанная выше (транслированный offset, консервативное округление вниз, LAG=0 после failover), проверена реально; полный сценарий «источник исчез навсегда, потом восстановился» — изложен здесь архитектурно, по документированному поведению MM2, но не подтверждён отдельным прогоном именно этого стенда.

RPO > 0 при асинхронной MM2-репликации — не недостаток конкретной настройки, а прямое следствие выбора между синхронным stretch-кластером и асинхронной репликацией отдельных кластеров, разобранного в начале статьи. Если бизнес-требование — нулевая потеря данных при аварии региона, ответ не «настроить MM2 быстрее», а «пересмотреть, нужен ли вообще синхронный stretch вместо гео-репликации», принимая на себя латентность межрегионального acks=all.

flowchart LR subgraph East["us-east (primary)"] P["Продюсеры/консьюмеры
(нормальный режим)"] T1["topic: orders"] P --> T1 end subgraph West["us-west (standby)"] T2["topic: us-east.orders
(префикс — от циклов)"] end MM2["MirrorMaker 2
Source + Checkpoint + Heartbeat"] T1 -->|"MirrorSourceConnector
асинхронно"| MM2 MM2 --> T2 T1 -.->|"MirrorCheckpointConnector
транслирует offset"| MM2 MM2 -.->|"sync.group.offsets"| T2 Fail(["Авария региона us-east"]) -.->|"failover"| West style East fill:#eee8d8,stroke:#8b7355 style West fill:#c9e4c5,stroke:#5b8a5e style MM2 fill:#f9f3e3,stroke:#8b7355 style Fail fill:#e8b4a0,stroke:#a8503a

flowchart LR
    subgraph East["us-east (primary)"]
        P["Продюсеры/консьюмеры
(нормальный режим)"] T1["topic: orders"] P --> T1 end subgraph West["us-west (standby)"] T2["topic: us-east.orders
(префикс — от циклов)"] end MM2["MirrorMaker 2
Source + Checkpoint + Heartbeat"] T1 -->|"MirrorSourceConnector
асинхронно"| MM2 MM2 --> T2 T1 -.->|"MirrorCheckpointConnector
транслирует offset"| MM2 MM2 -.->|"sync.group.offsets"| T2 Fail(["Авария региона us-east"]) -.->|"failover"| West style East fill:#eee8d8,stroke:#8b7355 style West fill:#c9e4c5,stroke:#5b8a5e style MM2 fill:#f9f3e3,stroke:#8b7355 style Fail fill:#e8b4a0,stroke:#a8503a
Active-passive DR: репликация, префикс, трансляция offset

Цена гео-репликации

Латентность. У синхронного stretch-кластера цена в пути записи — RTT до самой дальней реплики ISR входит в каждый acks=all-коммит напрямую. У MM2 локальная запись этой цены не несёт: producer на us-east не ждёт ничего с us-west, задержка проявляется только как лаг между кластерами, а не как задержка отдельной записи.

Трафик. Межрегиональный (тем более межконтинентальный) трафик у облачных провайдеров, как правило, платный и заметный по объёму — MM2 копирует полный поток реплицируемых топиков, и выбор, какие именно топики реплицировать (тот самый явный список topics = orders, а не .*), — не только вопрос порядка, но и вопрос счёта за трафик. Сжатие на стороне producer’ов (см. статью про хранение и compaction) снижает объём, который вообще нужно передавать между регионами.

Консистентность. Между кластерами, связанными через MM2, действует eventual-консистентность: гарантий уровня «эта запись видна на target сразу после подтверждения на source» здесь нет и не может быть — это структурное следствие асинхронности, а не настраиваемый параметр. Глобального порядка между кластерами и сквозного exactly-once через MM2 тоже нет: идемпотентность и транзакции (разобранные в статье про exactly-once semantics) действуют внутри одного кластера, MM2 их через границу кластеров не расширяет.

Что дальше

Здесь заканчивается серия о самой Kafka — от лога и партиций до гео-репликации между кластерами. Логичное продолжение — что делать с данными, когда они уже в Kafka: Kafka Streams и Flink разбирают stream processing поверх темы, которую эта серия закладывала с первой статьи (лог как источник, партиция как единица параллелизма). А если данные в Kafka должны попадать не из собственного кода, а прямо из базы данных, — статья про CDC через Debezium показывает связку PostgreSQL → Kafka через тот же Connect, на котором построен и MM2. Гео-репликация данных совсем другой природы (координаты, расстояния) — за пределами messaging — разобрана в статье про основы гео-поиска. Для навигации по всей теме messaging — карта messaging-landscape-map.

Обсуждение в Telegram

Присоединиться →

Комментарии