Бекапирование и восстановление с использование S3 хранилища#

Corax поддерживает сохранение и восстановление топиков в S3 хранилища. Функциональность реализуется с помощью механизма Kafka Connect и доработанного плагина Lenses Stream Reactor. В качестве примера S3 хранилища используется AWS S3 (данное хранилище используется в исходном плагине).

Для работы с S3 хранилищем:

  1. Настройте два коннектора: для записи данных в S3 (sink коннектор) и для получения данных из S3 (source коннектор).

  2. Задайте путь до плагина в параметре plugin.path, например: plugin.path={{ CONNECT_S3_PLUGIN_PATH }}.

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

Общая часть (для всех коннекторов)#

# Основные настройки Kafka
bootstrap.servers={host}:{port}  # Список брокеров Kafka в формате host:port (например: kafka1:9092, kafka2:9092)

# Конвертеры сообщений. Используются для шифрования/дешифрования ключей, значений и заголовков
# Укажите один и тот же тип конвертера для всех полей, если требуется единая политика шифрования
key.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter  # Конвертер для преобразования ключа
key.converter.password=<password>  # Пароль для дешифровки ключей
value.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Конвертер для преобразования значения
value.converter.password=<password>  # Пароль для дешифровки значений
header.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Конвертер, используемый для преобразования заголовков сообщений
header.converter.password=<password>  # Пароль для заголовков

# Включает передачу схемы вместе с данными (если используется формат, поддерживающий схемы)
key.converter.schemas.enable=true # Включение схемы для ключей
value.converter.schemas.enable=true  # Включение схемы для значений

# Идентификатор группы Kafka Connect
group.id=s3-connect-cluster  # Уникальное имя группы кластера коннекторов

# Топики для хранения состояния коннекторов
offset.storage.topic=connect-offsets # Топик, в котором хранятся смещения
offset.storage.replication.factor=1 # Фактор репликации для топика offset.storage.topic
offset.storage.partitions=1 # Количество партиций для топика offset.storage.topic
config.storage.topic=connect-configs # Топик, в котором хранится конфигурация
config.storage.replication.factor=1 # Фактор репликации
status.storage.topic=connect-status # Топик, в котором хранятся статусы выполняющихся задач (tasks) внутри коннекторов (connectors)
status.storage.replication.factor=1 # Фактор репликации для топика status.storage.topic
status.storage.partitions=1 # Количество партиций для топика status.storage.topic

# Интервал сброса смещений
offset.flush.interval.ms=1000

# Путь к плагину S3
plugin.path=/opt/corax-dev/optional/stream-reactor  # Директория с JAR-файлами коннектора

# Обработка ошибок
errors.log.enable=true  # Запись ошибок в лог
errors.log.include.messages=true  # Включение сообщений Kafka в логи ошибок
errors.deadletterqueue.topic.name=my-connector-errors  # Топик для сообщений, которые не удалось обработать (Dead Letter Queue)
errors.deadletterqueue.topic.replication.factor=1 # Фактор репликации для топика errors.deadletterqueue.topic.name
errors.tolerance=none  # Режим обработки ошибок (none/fail/all)

Sink Connector (запись в S3)#

# Идентификатор коннектора
name=local-s3-sink  # Уникальное имя коннектора

# Класс коннектора для S3
connector.class=io.lenses.streamreactor.connect.aws.s3.sink.S3SinkConnector

# Топики-источники данных
topics=test  # Список топиков Kafka для бекапа (через запятую)

# Настройки S3
connect.s3.custom.endpoint=https://{host}:{port}  # Кастомный endpoint S3
connect.s3.vhost.bucket=true  # Использовать виртуальный хостинг бакетов (true для AWS)
connect.s3.aws.auth.mode=Credentials  # Режим аутентификации (Credentials/IAM/Default)
connect.s3.aws.region=eu-west-2  # Регион AWS
connect.s3.aws.access.key=none  # Access Key для AWS
connect.s3.aws.secret.key=none  # Secret Key для AWS

# Правила экспорта данных (KCQL)
connect.s3.kcql=insert into test select * from test STOREAS `BYTES` PROPERTIES ('flush.count'=1)
# 'flush.count'=1: немедленная запись при получении сообщения

# Настройки безопасности Kafka (SSL)
consumer.ssl.trustStore.type=JKS  # Тип хранилища сертификатов
consumer.ssl.truststore.location=/mnt/security/secman.truststore.jks  # Путь к хранилищу доверенных сертификатов
consumer.ssl.truststore.password=<password>  # Пароль хранилища доверенных сертификатов
consumer.ssl.endpoint.identification.algorithm=  # Отключение проверки хоста
consumer.ssl.enabled.protocols=TLSv1.2  # Версия TLS
consumer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256  # Шифры, которые потребитель использует при установке TLS-соединения с брокерами
consumer.ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory  # Пользовательская SSL-фабрика

# Интеграция с {{names.secman.short}} (Vault)
consumer.secman.endpoint=https://secman-dzo.solution.sbt:8443  # URL {{names.secman.short}}
consumer.secman.namespace=DEV_DZO  # Пространство имен
consumer.secman.role_id=${decode:...}  # Role ID (закодировано)
consumer.secman.secret_id=${decode:...}  # Secret ID (закодировано)
consumer.secman.fetch.cn=<Common Name>  # Common Name для сертификата
consumer.secman.fetch.mount.path=PKI  # Точка монтирования PKI
consumer.secman.fetch.role=role-ga-secman-kafka-test  # Роль доступа

Source Connector (чтение из S3)#

# Идентификатор коннектора
name=local-s3-source  # Уникальное имя (не должно совпадать с sink!)

# Класс коннектора
connector.class=io.lenses.streamreactor.connect.aws.s3.source.S3SourceConnector

# Целевой топик Kafka
topic=test  # Топик для восстановленных данных

# Настройки S3 (аналогичны sink)
connect.s3.custom.endpoint=https://{host}:{port}
connect.s3.vhost.bucket=true
connect.s3.aws.auth.mode=Credentials
connect.s3.aws.region=eu-west-2
connect.s3.aws.access.key=none
connect.s3.aws.secret.key=none

# Правила импорта данных (KCQL)
connect.s3.kcql=insert into test select * from test STOREAS `BYTES`

# Настройки безопасности (SSL для producer)
producer.ssl.trustStore.type=JKS # Тип хранилища доверенных сертификатов
producer.ssl.truststore.location=/mnt/security/secman.truststore.jks # Путь к файлу хранилища доверенных сертификатов
producer.ssl.truststore.password=qwe123 # Пароль для доступа к хранилищу сертификатов
producer.ssl.endpoint.identification.algorithm=  # Метод идентификации конечной точки
producer.ssl.enabled.protocols=TLSv1.2 # Протокол, используемый для доступа к хранилищу сертификатов
producer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256 # Шифры, которые потребитель использует при установке TLS-соединения с брокерами
producer.ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory  # Пользовательская SSL-фабрика

# Интеграция с {{names.secman.short}}
producer.secman.endpoint=https://secman-dzo.solution.sbt:8443
producer.secman.namespace=DEV_DZO
producer.secman.role_id=${decode:...}
producer.secman.secret_id=${decode:...}
producer.secman.fetch.cn=<Common Name>
producer.secman.fetch.mount.path=PKI
producer.secman.fetch.role=role-ga-secman-kafka-test

Запуск#

Выполните команду:

bin/connect-distributed.sh connect.properties

Управление#

Управление осуществляется через REST API.