Ручная установка#
Последовательность действий#
Порядок ручной установки Corax из дистрибутива:
Внимание
Перед началом ручной установки необходимо собрать дистрибутив (zip-архив) по инструкции «Подготовка дистрибутива Corax к установке», если он не был собран ранее.
Извлеките дистрибутив из 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.shSSL с авторизацией через Secret Management System
deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth__audit__secman.shSSL с авторизацией через Secret Management System + аудит
deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth__audit.shSSL с авторизацией + аудит
deploy-ssl__zk_mtls_with_auth__kafka_ssl_with_auth.shSSL с авторизацией
deploy-plaintext__zk_plain_no_auth__kafka_plaintext_no_auth.shБез SSL и без авторизации
Перейдите в каталог
coraxкомандой:cd ~/coraxЗапустите скрипт с необходимым профилем, например:
./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.
Порядок настройки:
Установите переменную среды для домашнего каталога Corax (каталога, где установлен Corax). Например:
export CORAX_HOME=~/coraxДобавьте каталог
binвPATH:export PATH=$PATH:$CORAX_HOME/binПроверьте правильность установки переменной
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 выполните команды:
Для запуска ZooKeeper:
zookeeper-server-start -daemon ./config/zookeeper.properties;Для запуска брокера 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#
Перед запуском Schema Registry отредактируйте файл
${CORAX_HOME}/config/schemaregistry/schema-registry.properties:spring.kafka.bootstrap-servers=:9092где
9092— порт, на котором работает брокер Kafka.Для запуска Schema Registry используйте команду:
crx-schema-registry-start.sh -daemon config/schemaregistry/schema-registry.properties
Corax UI#
Перед запуском Corax UI отредактируйте файл
${CORAX_HOME}/config/kafka-ui.properties:kafka.clusters[0].bootstrapServers=:9092 kafka.clusters[0].zookeeper=:2181где
[0]— порядковый номер брокера (Corax UI может управлять несколькими брокерами или несколькими кластерами Kafka).Для запуска выполните команду:
crx-ui-start -daemon ./config/kafka-ui.properties
Запуск Kafka Connect#
Отредактируйте
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] #### #
Пример полного конфигурационного файла приведен в разеделе Пример настройки.
Запустите 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.
Проверка результата#
Проверка работоспособности на хостах группы Kafka#
Подключитесь по ssh к хосту из ansible группы kafka.
Выполните команду
ps axu | grep -v 'grep' | grep java.*kafka/server.properties. Должен отобразиться 1 запущенный процесс.Выполните пункты 1-2 для остальных хостов ansible группы kafka.
Проверка работоспособности на хостах группы ZooKeeper#
Подключитесь по ssh к хосту из ansible группы ZooKeeper.
Выполните команду
ps axu | grep -v 'grep' | grep java.*config/zookeeper.properties. Должен отобразиться 1 запущенный процесс.Выполните пункты 1-2 для остальных хостов ansible группы zookeeper.
Проверка работоспособности подключением через Kafka Tool / Offset Explorer#
Запустите Offset Explorer / Kafka Tool.
В элементе Cluster выберите Import Connection.
Импортируйте 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>
Подключитесь к кластеру и убедитесь, что на вкладке Brokers видны все брокеры, на которые выполнялась установка.
Проверка работоспособности UI#
Проверка работоспособности на хостах группы crxui:
Подключитесь по ssh к хосту из ansible группы crxui.
Выполните команду
ps axu | grep -v 'grep' | grep java.*kafka-ui.properties. Должен отобразиться 1 запущенный процесс.
Проверка работоспособности подключением через браузер:
Замените в адресе
http://ххх.ххх.ххх.хххх:pppp:Вместо
xxx.xxx.xxx.xxxукажите IP хоста из ansible группы crxui.Вместо
ppppукажите порт, установленный в параметреserver.portв файле/tmp/installer/ansible/inventories/DEV/group_vars/all/vars.yaml.
С помощью браузера перейдите по полученному адресу.
Проверка работоспособности Schema Registry#
Проверка работоспособности на хостах группы crxui:
Подключитесь по ssh к хосту из ansible группы crxui.
Выполните команду
ps axu | grep -v 'grep' | grep java.*schema-registry.properties. Должен отобразиться 1 запущенный процесс.
Проверка работоспособности подключением через браузер:
Замените в адресе
http://ххх.ххх.ххх.хххх:pppp:Вместо
xxx.xxx.xxx.xxxукажите IP хоста из ansible группы crxsr.Вместо
ppppукажите порт, установленный в параметреserver_portв файле/tmp/installer/ansible/inventories/DEV/group_vars/all/vars.yaml.
С помощью браузера перейдите по полученному адресу.
Проверка работоспособности Kafka Connect#
Проверьте отсутствие ошибок в логе
logs/kafka-connect.logВыполните запрос:
при работе под профилем
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:
Разверните кластер с настроенной интеграцией с компонентом COTE.
Проверьте, фиксируется ли событие от Corax в компонент COTE.
Проверка интеграции Secret Management System#
В случае успешной интеграции с Secret Management System в логах Corax должно появиться сообщение вида:
[2023-10-03 16:50:56,656] INFO Fetch secret keystore [cfg=SecmanConfig{<информация о конфигурации>}]