Kafka Connect для S3#
Отредактируйте
config/connect-s3.properties:Укажите адреса подключения к брокерам:
bootstrap.servers=10.ХХ.ХХ.ХХ:9093,xx.xx.xx.xx:9093,xx.xx.xx.xx:9093.Укажите адрес и порт, на котором будет работать Kafka Connect:
listeners=http://<IP адрес сервера Kafka Connect>:<порт>.Настройте работу с сертификатами:
Без сертификатов — дополнительных настроек не требуется, только выполнить указанные выше.
Пример
# Список брокеров Kafka в формате host:port (например, kafka1:9092,kafka2:9092) bootstrap.servers=10.ХХ.ХХ.ХХ:9093,10.ХХ.ХХ.ХХ:9093,10.ХХ.ХХ.ХХ:9093 listeners=https://10.ХХ.ХХ.ХХ:8084 # Идентификатор группы Kafka Connect # Уникальное имя группы кластера коннекторов group.id=s3-connect-cluster # Путь к плагину S3 # Директория с JAR-файлами коннектора plugin.path=/opt/Apache/kafka/optional/stream-reactor ############ COMMON ############### # Конвертеры сообщений. Используются для шифрования/дешифрования ключей, значений и заголовков # Укажите один и тот же тип конвертера для всех полей, если требуется единая политика шифрования # Конвертер для преобразования ключа key.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки ключей key.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования значения value.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки значений value.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования заголовков header.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки заголовков header.converter.password=${decode:<зашифрованное значение>} # Включает передачу схемы вместе с данными (если используется формат, поддерживающий схемы) # Включение схемы для ключей key.converter.schemas.enable=true # Включение схемы для значений value.converter.schemas.enable=true # Топики для хранения состояния коннекторов # Топик, в котором хранятся смещения offset.storage.topic=_crx-s3-connect-offsets # Фактор репликации для топика offset.storage.topic offset.storage.replication.factor=2 # Количество партиций для топика offset.storage.topic offset.storage.partitions=1 # Топики для хранения конфигурации коннекторов config.storage.topic=_crx-s3-connect-configs # Фактор репликации config.storage.replication.factor=2 # Топик, в котором хранятся статусы выполняющихся задач (tasks) внутри коннекторов (connectors) status.storage.topic=_crx-s3-connect-status # Фактор репликации для топика status.storage.topic status.storage.replication.factor=2 # Количество партиций для топика status.storage.topic status.storage.partitions=1 # Интервал сброса оффсетов #offset.flush.interval.ms=1000 # Обработка ошибок # Запись ошибок в лог errors.log.enable=true # Включение сообщений Kafka в логи ошибок errors.log.include.messages=true # Топик для сообщений, которые не удалось обработать (Dead Letter Queue) errors.deadletterqueue.topic.name=my-connector-errors # Фактор репликации для топика errors.deadletterqueue.topic.name errors.deadletterqueue.topic.replication.factor=1 # Режим обработки ошибок (none/fail/all) errors.tolerance=none config.providers=decode config.providers.decode.class=ru.sbt.ss.kafka.DecryptionConfigProvider config.providers.decode.param.security.encoding.class=ru.sbt.ss.password.decoder.SimpleTextPasswordDecoder config.providers.decode.param.security.encoding.salt=ru.sbt.ss.password.salt.SbtSaltProvider config.providers.decode.param.security.encoding.key=/opt/Apache/kafka/config/encrypt.passС SSL:
listeners.https.ssl.keystore.location=<путь к jks-файлу>— (для Kafka Connect REST API) расположение закрытого ключа в файле хранилища ключей;listeners.https.ssl.keystore.password=<пароль>— (для Kafka Connect REST API) пароль хранилища для файла хранилища ключей;listeners.https.ssl.key.password=<пароль>— (для Kafka Connect REST API) пароль приватного ключа в хранилище ключей;listeners.https.ssl.truststore.location=<путь к jks-файлу>— (для Kafka Connect REST API) путь к хранилищу доверенных сертификатов;listeners.https.ssl.truststore.password=<пароль>— (для Kafka Connect REST API) пароль к хранилищу доверенных сертификатов;ssl.keystore.location=<путь к jks-файлу>— (для брокера) расположение закрытого ключа в файле хранилища ключей;ssl.keystore.password=<пароль>— (для брокера) пароль хранилища для файла хранилища ключей;ssl.key.password=<пароль>— (для брокера) пароль приватного ключа в хранилище ключей;ssl.truststore.location=<путь к jks-файлу>— (для брокера) путь к хранилищу доверенных сертификатов;ssl.truststore.password=<пароль>— (для брокера) пароль к хранилищу доверенных сертификатов.
Сертификаты для подключения к Kafka:
producer.ssl.keystore.location=<путь к jks-файлу>— (для производителя) расположение закрытого ключа в файле хранилища ключей;producer.ssl.keystore.password=<пароль>— (для производителя) пароль хранилища для файла хранилища ключей;producer.ssl.key.password=<пароль>— (для производителя) пароль приватного ключа в хранилище ключей;producer.ssl.truststore.location=<путь к jks-файлу>— (для производителя) путь к хранилищу доверенных сертификатов;producer.ssl.truststore.password=<пароль>— (для производителя) пароль к хранилищу доверенных сертификатов;consumer.ssl.keystore.location=<путь к jks-файлу>— (для потребителя) расположение закрытого ключа в файле хранилища ключей;consumer.ssl.keystore.password=<пароль>— (для потребителя) пароль хранилища для файла хранилища ключей;consumer.ssl.key.password=<пароль>— (для потребителя) пароль приватного ключа в хранилище ключей;consumer.ssl.truststore.location=<путь к jks-файлу>— (для потребителя) путь к хранилищу доверенных сертификатов;consumer.ssl.truststore.password=<пароль>— (для потребителя) пароль к хранилищу доверенных сертификатов;
Пример
# Список брокеров Kafka в формате host:port (например, kafka1:9092,kafka2:9092) bootstrap.servers=10.ХХ.ХХ.ХХ:9093 listeners=https://10.ХХ.ХХ.ХХ:8083 # Идентификатор группы Kafka Connect # Уникальное имя группы кластера коннекторов group.id=s3-connect-cluster # Путь к плагину S3 # Директория с JAR-файлами коннектора plugin.path=/opt/Apache/kafka/optional/stream-reactor ############ COMMON ############### # Конвертеры сообщений. Используются для шифрования/дешифрования ключей, значений и заголовков # Укажите один и тот же тип конвертера для всех полей, если требуется единая политика шифрования # Конвертер для преобразования ключа key.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки ключей key.converter.password=change_it #key.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования значения value.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки значений value.converter.password=change_it #value.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования заголовков header.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки заголовков header.converter.password=change_it #header.converter.password=${decode:<зашифрованное значение>} # Включает передачу схемы вместе с данными (если используется формат, поддерживающий схемы) # Включение схемы для ключей key.converter.schemas.enable=true # Включение схемы для значений value.converter.schemas.enable=true # Топики для хранения состояния коннекторов # Топик, в котором хранятся смещения offset.storage.topic=_crx-s3-connect-offsets # Фактор репликации для топика offset.storage.topic offset.storage.replication.factor=2 # Количество партиций для топика offset.storage.topic offset.storage.partitions=1 # Топики для хранения конфигурации коннекторов config.storage.topic=_crx-s3-connect-configs # Фактор репликации config.storage.replication.factor=2 # Топик, в котором хранятся статусы выполняющихся задач (tasks) внутри коннекторов (connectors) status.storage.topic=_crx-s3-connect-status # Фактор репликации для топика status.storage.topic status.storage.replication.factor=2 # Количество партиций для топика status.storage.topic status.storage.partitions=1 # Интервал сброса оффсетов #offset.flush.interval.ms=1000 # Обработка ошибок # Запись ошибок в лог errors.log.enable=true # Включение сообщений Kafka в логи ошибок errors.log.include.messages=true # Топик для сообщений, которые не удалось обработать (Dead Letter Queue) errors.deadletterqueue.topic.name=my-connector-errors # Фактор репликации для топика errors.deadletterqueue.topic.name errors.deadletterqueue.topic.replication.factor=1 # Режим обработки ошибок (none/fail/all) errors.tolerance=none # ===================== SSL params =============== listeners.https.ssl.keystore.location=ssl/esbmon.jks listeners.https.ssl.keystore.password=change_it #listeners.https.ssl.keystore.password=${decode:<зашифрованное значение>} listeners.https.ssl.key.password=change_it #listeners.https.ssl.key.password=${decode:<зашифрованное значение>} listeners.https.ssl.keystore.type=JKS listeners.https.ssl.truststore.type=JKS listeners.https.ssl.truststore.location=ssl/esbmon.jks listeners.https.ssl.truststore.password=change_it #listeners.https.ssl.truststore.password=${decode:<зашифрованное значение>} listeners.https.ssl.endpoint.identification.algorithm= listeners.https.ssl.enabled.protocols=TLSv1.2 listeners.https.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 listeners.https.ssl.client.auth=required security.protocol=SSL ssl.keystore.location=ssl/esbmon.jks ssl.keystore.password=change_it #ssl.keystore.password=${decode:<зашифрованное значение>} ssl.key.password=change_it #ssl.key.password=${decode:<зашифрованное значение>} ssl.keystore.type=JKS ssl.truststore.type=JKS ssl.truststore.location=ssl/esbmon.jks ssl.truststore.password=change_it #ssl.truststore.password=${decode:<зашифрованное значение>} ssl.endpoint.identification.algorithm= ssl.enabled.protocols=TLSv1.2 ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 ssl.client.auth=required producer.security.protocol=SSL producer.ssl.keystore.location=ssl/esbmon.jks producer.ssl.keystore.password=change_it #producer.ssl.keystore.password=${decode:<зашифрованное значение>} producer.ssl.key.password=change_it #producer.ssl.key.password=${decode:<зашифрованное значение>} producer.ssl.keystore.type=JKS producer.ssl.truststore.type=JKS producer.ssl.truststore.location=ssl/esbmon.jks producer.ssl.truststore.password=change_it #producer.ssl.truststore.password=${decode:<зашифрованное значение>} producer.ssl.endpoint.identification.algorithm= producer.ssl.enabled.protocols=TLSv1.2 producer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 consumer.security.protocol=SSL consumer.ssl.keystore.location=ssl/esbmon.jks consumer.ssl.keystore.password=change_it #consumer.ssl.keystore.password=${decode:<зашифрованное значение>} consumer.ssl.key.password=change_it #consumer.ssl.key.password=${decode:<зашифрованное значение>} consumer.ssl.keystore.type=JKS consumer.ssl.truststore.type=JKS consumer.ssl.truststore.location=ssl/esbmon.jks consumer.ssl.truststore.password=change_it #consumer.ssl.truststore.password=${decode:<зашифрованное значение>} consumer.ssl.endpoint.identification.algorithm= consumer.ssl.enabled.protocols=TLSv1.2 consumer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 config.providers=decode config.providers.decode.class=ru.sbrf.kafka.config.provider.DecryptionConfigProvider config.providers.decode.param.security.encoding.class=ru.sbt.ss.password.decoder.SimpleTextPasswordDecoder config.providers.decode.param.security.encoding.salt=ru.sbt.ss.password.salt.SbtSaltProvider config.providers.decode.param.security.encoding.key=/opt/Apache/kafka/config/encrypt.pass ##### [CUSTOM_PROPERTIES] #### #C сертификатами с получением через SecMan (все параметры
secman.*):secman.endpoint=<адрес Sec_Man>:<порт>— URL SecMan;secman.namespace=<пространство имен>— пространство имен;secman.role_id=<закодированное значение>— Role ID (закодировано);secman.secret_id=<закодированное значение>— Secret ID (закодировано);secman.fetch.cn=<CN>— Common Name для сертификата;secman.fetch.mount.path=<PKI>— точка монтирования PKI;secman.fetch.role=<роль>— роль для доступа.
Пример
# Список брокеров Kafka в формате host:port (например, kafka1:9092,kafka2:9092) bootstrap.servers=10.ХХ.ХХ.ХХ:9093 listeners=https://10.ХХ.ХХ.ХХ:8083 # Идентификатор группы Kafka Connect # Уникальное имя группы кластера коннекторов group.id=s3-connect-cluster # Путь к плагину S3 # Директория с JAR-файлами коннектора plugin.path=/opt/Apache/kafka/optional/stream-reactor ############ COMMON ############### # Конвертеры сообщений. Используются для шифрования/дешифрования ключей, значений и заголовков # Укажите один и тот же тип конвертера для всех полей, если требуется единая политика шифрования # Конвертер для преобразования ключа key.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки ключей key.converter.password=change_it #key.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования значения value.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки значений value.converter.password=change_it #value.converter.password=${decode:<зашифрованное значение>} # Конвертер для преобразования заголовков header.converter=ru.sbrf.kafka.connect.EncryptedByteArrayConverter # Пароль для дешифровки заголовков header.converter.password=change_it #header.converter.password=${decode:<зашифрованное значение>} # Включает передачу схемы вместе с данными (если используется формат, поддерживающий схемы) # Включение схемы для ключей key.converter.schemas.enable=true # Включение схемы для значений value.converter.schemas.enable=true # Топики для хранения состояния коннекторов # Топик, в котором хранятся смещения offset.storage.topic=_crx-s3-connect-offsets # Фактор репликации для топика offset.storage.topic offset.storage.replication.factor=2 # Количество партиций для топика offset.storage.topic offset.storage.partitions=1 # Топики для хранения конфигурации коннекторов config.storage.topic=_crx-s3-connect-configs # Фактор репликации config.storage.replication.factor=2 # Топик, в котором хранятся статусы выполняющихся задач (tasks) внутри коннекторов (connectors) status.storage.topic=_crx-s3-connect-status # Фактор репликации для топика status.storage.topic status.storage.replication.factor=2 # Количество партиций для топика status.storage.topic status.storage.partitions=1 # Интервал сброса оффсетов #offset.flush.interval.ms=1000 # Обработка ошибок # Запись ошибок в лог errors.log.enable=true # Включение сообщений Kafka в логи ошибок errors.log.include.messages=true # Топик для сообщений, которые не удалось обработать (Dead Letter Queue) errors.deadletterqueue.topic.name=my-connector-errors # Фактор репликации для топика errors.deadletterqueue.topic.name errors.deadletterqueue.topic.replication.factor=1 # Режим обработки ошибок (none/fail/all) errors.tolerance=none # ===================== SSL params =============== listeners.https.ssl.trustStore.type=JKS listeners.https.ssl.truststore.location=ssl/secman-truststore.jks listeners.https.ssl.truststore.password=change_it #listeners.https.ssl.truststore.password=${decode:<зашифрованное значение>} listeners.https.ssl.endpoint.identification.algorithm= listeners.https.ssl.enabled.protocols=TLSv1.2 listeners.https.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 listeners.https.ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory listeners.https.secman.endpoint=<адрес SecMan>:<порт> listeners.https.secman.namespace=<пространство имен> listeners.https.secman.role_id=change_it #listeners.https.secman.role_id=${decode:<зашифрованное значение>} listeners.https.secman.secret_id=change_it #listeners.https.secman.secret_id=${decode:<зашифрованное значение>} listeners.https.secman.fetch.cn=<CN> listeners.https.secman.fetch.mount.path=<PKI> listeners.https.secman.fetch.role=<роль для доступа> listeners.https.ssl.client.auth=required security.protocol=SSL ssl.trustStore.type=JKS ssl.truststore.location=ssl/secman-truststore.jks ssl.truststore.password=change_it #ssl.truststore.password=${decode:<зашифрованное значение>} ssl.endpoint.identification.algorithm= ssl.enabled.protocols=TLSv1.2 ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory ssl.client.auth=required secman.endpoint=<адрес SecMan>:<порт> secman.namespace=<пространство имен> secman.role_id=change_it #secman.role_id=${decode:<зашифрованное значение>} secman.secret_id=change_it #secman.secret_id=${decode:<зашифрованное значение>} secman.fetch.cn=<CN> secman.fetch.mount.path=<PKI> secman.fetch.role=<роль для доступа> producer.security.protocol=SSL producer.ssl.trustStore.type=JKS producer.ssl.truststore.location=ssl/secman-truststore.jks producer.ssl.truststore.password=change_it #producer.ssl.truststore.password=${decode:<зашифрованное значение>} producer.ssl.endpoint.identification.algorithm= producer.ssl.enabled.protocols=TLSv1.2 producer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 producer.ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory producer.secman.endpoint=<адрес SecMan>:<порт> producer.secman.namespace=<пространство имен> producer.secman.role_id=change_it #producer.secman.role_id=${decode:<зашифрованное значение>} producer.secman.secret_id=change_it #producer.secman.secret_id=${decode:<зашифрованное значение>} producer.secman.fetch.cn=<CN> producer.secman.fetch.mount.path=<PKI> producer.secman.fetch.role=<роль для доступа> consumer.security.protocol=SSL consumer.ssl.trustStore.type=JKS consumer.ssl.truststore.location=ssl/secman-truststore.jks consumer.ssl.truststore.password=change_it #consumer.ssl.truststore.password=${decode:<зашифрованное значение>} consumer.ssl.endpoint.identification.algorithm= consumer.ssl.enabled.protocols=TLSv1.2 consumer.ssl.cipher.suites=TLS_ECDHE_RSA_WITH_AES_128_GCM_SHA256,TLS_ECDHE_ECDSA_WITH_AES_128_GCM_SHA256 consumer.ssl.engine.factory.class=ru.sbrf.kafka.secman.SecmanSslEngineFactory consumer.secman.endpoint=<адрес SecMan>:<порт> consumer.secman.namespace=<пространство имен> consumer.secman.role_id=change_it #consumer.secman.role_id=${decode:<зашифрованное значение>} consumer.secman.secret_id=change_it #consumer.secman.secret_id=${decode:<зашифрованное значение>} consumer.secman.fetch.cn=<CN> consumer.secman.fetch.mount.path=<PKI> consumer.secman.fetch.role=<роль для доступа> config.providers=decode config.providers.decode.class=ru.sbrf.kafka.config.provider.DecryptionConfigProvider config.providers.decode.param.security.encoding.class=ru.sbt.ss.password.decoder.SimpleTextPasswordDecoder config.providers.decode.param.security.encoding.salt=ru.sbt.ss.password.salt.SbtSaltProvider config.providers.decode.param.security.encoding.key=/opt/Apache/kafka/config/encrypt.pass ##### [CUSTOM_PROPERTIES] #### #
Пример полного конфигурационного файла приведен в разделе Пример настройки.
Запустите Kafka Connect из
KAFKA_HOMEдиректории:bin/connect-distributed.sh -daemon config/connect-s3.properties.
Создание и управление коннекторами для S3#
Подготовьте json-файлы с параметрами коннекторов. Параметры описаны в разделах:
Примеры JSON:
{ "config": { "connect.s3.aws.access.key": "none", "connect.s3.aws.auth.mode": "Credentials", "connect.s3.aws.region": "eu-west-2", "connect.s3.aws.secret.key": "none", "connect.s3.custom.endpoint": "http://10.xx.xx.xx:8082", "connect.s3.kcql": "insert into kafka select * from `*` STOREAS `JSON` PROPERTIES ('flush.count'=1,'store.envelope'=true)", "connect.s3.vhost.bucket": "true", "connector.class": "io.lenses.streamreactor.connect.aws.s3.sink.S3SinkConnector", "initial_state": "STOPPED", "name": "s3-sink-encrypt", "tasks.max": "1", "topics": "encryptSink,jsonSink" }, "initial_state": "STOPPED", "name": "s3-sink-encrypt" }{ "config": { "connect.s3.aws.access.key": "none", "connect.s3.aws.auth.mode": "Credentials", "connect.s3.aws.region": "eu-west-2", "connect.s3.aws.secret.key": "none", "connect.s3.custom.endpoint": "http://10.xx.xx.xx:8082", "connect.s3.kcql": "insert into encryptSource select * from kafka:encryptSink STOREAS `JSON` PROPERTIES ('store.envelope'=true); insert into jsonSource select * from kafka:jsonSink STOREAS `JSON` PROPERTIES ('store.envelope'=true)", "connect.s3.vhost.bucket": "true", "connector.class": "io.lenses.streamreactor.connect.aws.s3.source.S3SourceConnector", "initial_state": "STOPPED", "name": "s3-source-encrypt", "tasks.max": "1", "topics": "encryptSource,jsonSource" }, "initial_state": "STOPPED", "name": "s3-source-encrypt" }Создайте коннекторы отправив запрос в Kafka Connect:
curl -s --cert cert.pem --key cert.key -kX "POST" --header "Content-Type: application/json" --data @<путь к json-файлу, созданному на 1 шаге> "https://10.xx.xx.xx:8083/connectors/"Работа без сертификатов
При работе без сертификатов не используйте параметры
--certи--key.