Установка и запуск#

Для установки Corax Mirror Maker 2 (MM2) выполните следующие шаги:

  1. Убедитесь, что установлена Java в соответствии с системными требованиями:

    java -version
    
  2. Подготовьте файл конфигурации *mirror-maker.properties: создайте потоки репликации, настройте мониторинг и журналирование. Полный перечень возможных настроек приведен в разделе «Общие конфигурации Corax Mirror Maker 2».

  3. Запустите Corax Mirror Maker 2 командой:

    bin/connect-mirror-maker.sh -daemon config/connect-mirror-maker.properties
    

Остановка репликации#

Для остановки репликации отправьте SIGTERM сигнал, например командой:

$ kill <MirrorMaker pid>

Применение изменений настроек#

Для применения изменений в конфиг-файле перезапустите процесс Corax Mirror Maker 2.

Настройка#

Для настройки MM2 необходим файл конфигурации mirror-maker.properties, подробное описание параметров приведено в разделе «Общие конфигурации Corax Mirror Maker 2».

Настройки можно разделить на две группы:

Основные параметры#

Параметр

Описание

Пример

Значение по умолчанию

clusters

Список псевдонимов кластеров, через запятую (например, source, destination)

clusters = primary, backup

{cluster}.bootstrap.servers

Адреса брокеров для указанного кластера

primary.bootstrap.servers = kafka1:9092

{cluster}.alias

Алиас кластера (если имя отличается от идентификатора)

primary.alias = main-cluster

<src>-><dst>.enabled

Включает репликацию из кластера <src> в кластер <dst>

<src>-><dst>.topics

Регулярное выражение для отбора топиков для репликации

.*

replication.factor

Фактор репликации для вновь созданных топиков

2

groups

Перечень групп потребителей, подлежащих репликации

.*

topics.blacklist

Черный список топиков, исключаемых из репликации

null

tasks.max

Максимальное количество параллельных задач

1

Параметры безопасности#

Для каждого кластера:

Параметр

Описание

Пример

{cluster}.security.protocol

Протокол, используемый для связи с брокерами

primary.security.protocol = SSL

{cluster}.ssl.truststore.location

Местоположение хранилища доверенных сертификатов

primary.ssl.truststore.location = /certs/truststore.jks

ssl.truststore.password

Пароль для файла хранилища доверенных сертификатов

{cluster}.ssl.keystore.location

Расположение файла хранилища ключей

primary.ssl.keystore.location = /certs/keystore.jks

{cluster}.sasl.mechanism

Механизм SASL, используемый для подключения клиентов (например, PLAIN, SCRAM-SHA-512)

primary.sasl.mechanism = SCRAM-SHA-512

{cluster}.sasl.jaas.config

Контекстные параметры входа JAAS для SASL-подключений в формате конфигурационных файлов JAAS

primary.sasl.jaas.config = org.apache.kafka.common.security.scram.ScramLoginModule...

Расширенные параметры#

Параметр

Описание

kafka_mirrormaker.emit_checkpoints_enabled

Включает или отключает передачу контрольных точек смещения

whitelist

Белый список топиков для репликации

auto.create.topics.enable

Авто-регистрация отсутствующих топиков в целевом кластере

retention.ms

Срок хранения сообщений в миллисекундах

compression.type

Тип сжатия данных

Репликация между кластерами#

Для каждого направления репликации (source->target):

Параметр

Описание

Пример

Значение по умолчанию

{source}->{target}.enabled

Включить репликацию из source в target

primary->backup.enabled = true

false

{source}->{target}.topics

Регулярные выражения для включения топиков

primary->backup.topics = orders_.*

.* (все топики)

{source}->{target}.topics.exclude

Исключить топики из репликации

primary->backup.topics.exclude = audit_log

{source}->{target}.groups

Регулярные выражения для групп потребителей

primary->backup.groups = app_.*

.*

{source}->{target}.groups.exclude

Исключить группы потребителей

primary->backup.groups.exclude = test_group

Репликация данных#

Параметр

Описание

Пример

Значение по умолчанию

replication.factor

Фактор репликации для создаваемых топиков

replication.factor = 3

2

offset-syncs.topic.replication.factor

Фактор репликации для топика offset-syncs

offset-syncs.topic.replication.factor = 3

3

heartbeats.topic.replication.factor

Фактор репликации для топика heartbeats

heartbeats.topic.replication.factor = 3

3

checkpoints.topic.replication.factor

Фактор репликации для топика контрольных точек (checkpoints)

checkpoints.topic.replication.factor = 3

3

offset.storage.replication.factor

Фактор репликации, используемый при создании топика для хранения смещений (offsets) коннекторов и задач в Kafka Connect

offset.storage.replication.factor = 3

3

status.storage.replication.factor

Фактор репликации, используемый при создании топика для хранения статусов коннекторов и задач в Kafka Connect

status.storage.replication.factor = 3

3

config.storage.replication.factor

Фактор репликации, используемый при создании топика для хранения конфигураций коннекторов и задач в Kafka Connect

config.storage.replication.factor = 3

3

emit.checkpoints.interval.seconds

Частота отправки контрольных точек (в секундах)

emit.checkpoints.interval.seconds = 60

60

emit.heartbeats.interval.seconds

Частота отправки heartbeat-сигналов (в секундах)

emit.heartbeats.interval.seconds = 5

5

sync.topic.configs.enabled

Синхронизировать ли конфигурации создаваемых топиков с их исходными аналогами

sync.topic.configs.enabled = true

true

sync.topic.configs.interval.seconds

Интервал (в секундах) синхронизации конфигураций топиков

sync.topic.configs.interval.seconds = 300

300

sync.group.offsets.enabled

Записывать ли периодически преобразованные смещения в топик __consumer_offsets целевого кластера

sync.group.offsets.enabled = true

false

refresh.topics.interval.seconds

Интервал (в секундах) проверки новых топиков

refresh.topics.interval.seconds = 60

60

refresh.groups.interval.seconds

Интервал (в секундах) проверки новых групп потребителей

refresh.groups.interval.seconds = 60

60

Производительность и тюнинг#

Параметр

Описание

Пример

Значение по умолчанию

tasks.max

Максимальное количество задач для коннектора

tasks.max = 4

1

max.poll.records

Максимальное количество записей, возвращаемых за один вызов poll()

max.poll.records = 1000

500

replication.policy.class

Класс, определяющий соглашение об именовании создаваемых на целевом кластере топиков

replication.policy.class = com.custom.ReplicationPolicy

DefaultReplicationPolicy

offset.lag.max

Максимальное отставание (в количестве записей) удаленной партиции перед повторной синхронизацией

offset.lag.max = 1000

100

Мониторинг и логирование#

Параметр

Описание

Пример

metrics.reporter

Список классов, реализующих интерфейс org.apache.kafka.common.metrics.MetricsReporter. Эти классы используются для вывода метрик. Можно указать несколько репортеров, разделив их запятыми

metrics.reporter = org.apache.kafka.common.metrics.JmxReporter

metrics.recording.level

Уровень детализации метрик (DEBUG, INFO)

metrics.recording.level = INFO

client.id

Идентификатор клиента, который передается серверу при каждом запросе. Используется для логирования и мониторинга на стороне брокера

client.id = mm2-primary-consumer

Настройки Kafka Connect#

MM2 работает поверх Kafka Connect, поэтому поддерживает все его параметры:

Параметр

Описание

Пример

group.id

Идентификатор группы Kafka Connect

group.id = mm2-cluster

key.converter

Класс-конвертер для преобразования данных между форматом Corax Connect и сериализованной формой, записываемой в Corax

key.converter = org.apache.kafka.connect.json.JsonConverter

value.converter

Класс-конвертер для преобразования данных между форматом Corax Connect и сериализованной формой, записываемой в Kafka

value.converter = org.apache.kafka.connect.json.JsonConverter

offset.storage.topic

Группа для хранения смещений коннектора исходного топика

offset.storage.topic = mm2-offsets