Kafka Connect для S3#

  1. Отредактируйте 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] ####
        #
        
        

    Пример полного конфигурационного файла приведен в разделе Пример настройки.

  2. Запустите Kafka Connect из KAFKA_HOME директории: bin/connect-distributed.sh -daemon config/connect-s3.properties.

Создание и управление коннекторами для S3#

  1. Подготовьте 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"
    }
    
  2. Создайте коннекторы отправив запрос в 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.