Ручная установка#

Последовательность действий#

Порядок ручной установки Corax из дистрибутива:

Внимание

Перед началом ручной установки необходимо собрать дистрибутив (zip-архив) по инструкции «Подготовка дистрибутива Corax к установке», если он не был собран ранее.

  1. Извлеките дистрибутив из zip-архива командой:

    unzip KFK-10.340.0-16-distrib.zip -d ~/corax
    

    где KFK-10.340.0-16-distrib.zip — имя дистрибутива.

    Примечание

    В примере использовано имя дистрибутива Corax версии 10.340.0. Архив будет извлечен в каталог corax.

    В результате распаковки в каталоге corax должны появится файлы:

    Файл

    Описание

    Архивы

    kfka-doc-10.340.0-16-distrib.zip

    Архив с документацией

    kfka-deploy-10.340.0-SNAPSHOT-distrib.zip

    Архив с ansible-ролями

    kafka-sbt-config.zip

    Архив с конфигурациями для профилей установки

    kafka-dist.zip

    Дистрибутив Corax

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

    deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth__secman.sh

    SSL с авторизацией через Secret Management System

    deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth__audit__secman.sh

    SSL с авторизацией через Secret Management System + аудит

    deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth__audit.sh

    SSL с авторизацией + аудит

    deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth.sh

    SSL с авторизацией

    deploy-plaintext__zk_plain_no_auth__kafka_plaintext_no_auth.sh

    Без SSL и без авторизации

  2. Перейдите в каталог corax командой:

    cd ~/corax
    
  3. Запустите скрипт с необходимым профилем, например:

    ./deploy-plaintext__zk_plain_no_auth__kafka_plaintext_no_auth.sh
    

    Скрипт распакует дистрибутив, переложит файлы properties для выбранного профиля в каталог ./config. В результате в текущем каталоге должны появится каталоги:

    Каталог

    Описание

    bin

    Каталог с исполняемыми файлами

    config

    Каталог с конфигурационными файлами

    data

    Каталог для данных

    etc

    Каталог для символической ссылки kafka

    libs

    Каталог для библиотек kafka

    logs

    Каталог для логов kafka

    optional

    Каталог для опциональных инструментов

    util

    Каталог, в который перекладывается Corax после применения скрипта профиля

Настройка CORAX_HOME и PATH (обязательный)#

Чтобы упростить использование Corax CLI и всех инструментов командной строки, поставляемых вместе с Corax , настройте переменную CORAX_HOME и добавьте каталог bin в PATH. При этом станет возможным использование инструментов CLI без перехода в каталог CORAX_HOME.

Порядок настройки:

  1. Установите переменную среды для домашнего каталога Corax (каталога, где установлен Corax). Например:

    export CORAX_HOME=~/corax
    
  2. Добавьте каталог bin в PATH:

    export PATH=$PATH:$CORAX_HOME/bin
    
  3. Проверьте правильность установки переменной CORAX_HOME командой:

    kafka-configs.sh --version
    

    В выводе должны быть показаны Дата сборки, Версия Corax и базовая версия Corax, например:

    Core version: Х.Х.Х, Corax version: ХХ.ХХХ.Х-ХХ (Commit:<хеш значение>), build: DD/MM/YYYY
    

Запуск Corax#

Брокер Kafka и ZooKeeper#

Для запуска Corax отредактируйте файл ${CORAX_HOME}/config/server.properties:

zookeeper.connect=:2181
listeners=PLAINTEXT://:9092

Если предварительно была настроена переменная CORAX_HOME, то для запуска Corax выполните команды:

  1. Для запуска ZooKeeper: zookeeper-server-start -daemon ./config/zookeeper.properties;

  2. Для запуска брокера Kafka: kafka-server-start -daemon ./config/server.properties.

Если переменная CORAX_HOME не настроена, используйте для запуска Corax команды:

cd ~/corax;
bin/zookeeper-server-start.sh -daemon config/zookeeper.properties
bin/kafka-server-start -daemon config/server.properties

Schema Registry#

  1. Перед запуском Schema Registry отредактируйте файл ${CORAX_HOME}/config/schemaregistry/schema-registry.properties:

    spring.kafka.bootstrap-servers=:9092
    

    где 9092 — порт, на котором работает брокер Kafka.

  2. Для запуска Schema Registry используйте команду:

    crx-schema-registry-start.sh -daemon config/schemaregistry/schema-registry.properties
    

Corax UI#

  1. Перед запуском Corax UI отредактируйте файл ${CORAX_HOME}/config/kafka-ui.properties:

    kafka.clusters[0].bootstrapServers=:9092
    kafka.clusters[0].zookeeper=:2181
    

    где [0] — порядковый номер брокера (Corax UI может управлять несколькими брокерами или несколькими кластерами Kafka).

  2. Для запуска выполните команду:

    crx-ui-start -daemon ./config/kafka-ui.properties
    

Запуск Kafka Connect#

  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=<адрес SecMan>:<порт> — 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.

Проверка результата#

Проверка работоспособности на хостах группы Kafka#

  1. Подключитесь по ssh к хосту из ansible группы kafka.

  2. Выполните команду ps axu | grep -v 'grep' | grep java.*kafka/server.properties. Должен отобразиться 1 запущенный процесс.

  3. Выполните пункты 1-2 для остальных хостов ansible группы kafka.

Проверка работоспособности на хостах группы ZooKeeper#

  1. Подключитесь по ssh к хосту из ansible группы ZooKeeper.

  2. Выполните команду ps axu | grep -v 'grep' | grep java.*config/zookeeper.properties. Должен отобразиться 1 запущенный процесс.

  3. Выполните пункты 1-2 для остальных хостов ansible группы zookeeper.

Проверка работоспособности подключением через Kafka Tool / Offset Explorer#

  1. Запустите Offset Explorer / Kafka Tool.

  2. В элементе Cluster выберите Import Connection.

  3. Импортируйте xml файл с конфигурацией подключения из примеров ниже (при необходимости отредактируйте параметры подключения под свое окружение):

    • для протокола безопасности SSL__ZK_mTLS_WITH_AUTH__KAFKA_SSL_WITH_AUTH:

      <?xml version="1.0" encoding="UTF-8" standalone="no"?>
      <connections>
      <connection bootstrap_servers="<host-ip-address>:9093" broker_security_type="SSL" chroot="/" group="Clusters" groupId="1" host="10.XX.XX.XX" jaas_config="" keystore_location="C:\Users\User\Desktop\server.kafka.jks" keystore_password="<password>" keystore_privatekey="<password>" name="10.XX.XX.XX" port="2181" sasl_mechanism="" schema_registry_endpoint="" truststore_location="C:\Users\User\Desktop\truststore.kafka.jks" truststore_password="<password>" version="VERSION_1_0_0"/>
      <groups>
      <group id="1" name="Clusters"/>
      </groups>
      </connections>
      
    • для протокола безопасности PLAINTEXT__ZK_PLAIN_NO_AUTH__KAFKA_PLAINTEXT_NO_AUTH:

      <?xml version="1.0" encoding="UTF-8" standalone="no"?>
      <connections>
      <connection bootstrap_servers="<host-ip-address>:9093" broker_security_type="PLAINTEXT" chroot="/" group="Clusters" groupId="1" host="10.XX.XX.XX" jaas_config="" keystore_location="" keystore_password="" keystore_privatekey="" name="10.XX.XX.XX" port="2181" sasl_mechanism="" schema_registry_endpoint="" truststore_location="" truststore_password="" version="VERSION_1_0_0"/>
      <groups>
      <group id="1" name="Clusters"/>
      </groups>
      </connections>
      
  4. Подключитесь к кластеру и убедитесь, что на вкладке Brokers видны все брокеры, на которые выполнялась установка.

Проверка работоспособности UI#

Проверка работоспособности на хостах группы crxui:

  1. Подключитесь по ssh к хосту из ansible группы crxui.

  2. Выполните команду ps axu | grep -v 'grep' | grep java.*kafka-ui.properties. Должен отобразиться 1 запущенный процесс.

Проверка работоспособности подключением через браузер:

  1. Замените в адресе http://ххх.ххх.ххх.хххх:pppp:

    1. Вместо xxx.xxx.xxx.xxx укажите IP хоста из ansible группы crxui.

    2. Вместо pppp укажите порт, установленный в параметре server.port в файле /tmp/installer/ansible/inventories/DEV/group_vars/all/vars.yaml.

  2. С помощью браузера перейдите по полученному адресу.

Проверка работоспособности Schema Registry#

Проверка работоспособности на хостах группы crxui:

  1. Подключитесь по ssh к хосту из ansible группы crxui.

  2. Выполните команду ps axu | grep -v 'grep' | grep java.*schema-registry.properties. Должен отобразиться 1 запущенный процесс.

Проверка работоспособности подключением через браузер:

  1. Замените в адресе http://ххх.ххх.ххх.хххх:pppp:

    1. Вместо xxx.xxx.xxx.xxx укажите IP хоста из ansible группы crxsr.

    2. Вместо pppp укажите порт, установленный в параметре server_port в файле /tmp/installer/ansible/inventories/DEV/group_vars/all/vars.yaml.

  2. С помощью браузера перейдите по полученному адресу.

Проверка работоспособности Kafka Connect#

  1. Проверьте отсутствие ошибок в логе logs/kafka-connect.log

  2. Выполните запрос:

    • при работе под профилем PLAINTEXT (работа без сертификата): curl -s -kX "GET" --header "Content-Type: application/json" "http://xx.xx.xx.xx:8083"

    • при работе с остальными профилями — в запросе укажите сертификат (cert.pem и cert.key): curl -s --cert cert.pem --key cert.key -kX "GET" --header "Content-Type: application/json" "https://xx.xx.xx.xx:8083".

    ожидаемый ответ — информация о версии Corax, например: {"version":"3.9.1","commit":"a8b0f70c446943c9","kafka_cluster_id":"kS2HdtRiR4yzW9N0NlSQIA"}.

Проверка интеграции компонентом «Единый коллектор телеметрии» (COTE) продукта Platform V Monitor#

Для проверки настройки соединения между Corax и компонентом COTE:

  1. Разверните кластер с настроенной интеграцией с компонентом COTE.

  2. Проверьте, фиксируется ли событие от Corax в компонент COTE.

Проверка интеграции Secret Management System#

В случае успешной интеграции с Secret Management System в логах Corax должно появиться сообщение вида:

[2023-10-03 16:50:56,656] INFO Fetch secret keystore [cfg=SecmanConfig{<информация о конфигурации>}]