Бекапирование и восстановление с использование S3 хранилища#
Corax поддерживает сохранение и восстановление топиков в S3 хранилища. Функциональность реализуется с помощью механизма Kafka Connect и доработанного плагина Lenses Stream Reactor. В качестве примера S3 хранилища используется AWS S3 (данное хранилище используется в исходном плагине).
Для работы с S3 хранилищем:
Настройте два коннектора: для записи данных в S3 (sink коннектор) и для получения данных из S3 (source коннектор).
Задайте путь до плагина в параметре
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.