Добавление нового коннектора для Kafka Connect
Обзор
По умолчанию для создания в ADS доступны следующие коннекторы Kafka Connect:
Плагины этих коннекторов отображаются при создании коннектора в пользовательском интерфейсе ADS Control и размещаются в директории /usr/lib/kafka-connect/plugins, указанной в качестве значения параметра plugin.path в группе connect-distributed.properties на странице конфигурирования сервиса Kafka Connect.
Остальные коннекторы могут быть добавлены самостоятельно.
Добавление плагина коннектора
-
Подготовьте один из вариантов плагина коннектора:
-
исполняемый JAR-файл (uber-JAR), содержащий Java-код коннектора и все его зависимости;
-
поддиректория плагина со всеми JAR-файлами и зависимостями (рекомендуется одна поддиректория на каждый плагин).
ВНИМАНИЕПлагин не должен включать runtime-библиотеки Kafka Connect.
Для примера в данной статье используется исполняемый JAR-файл плагина для ClickHouseSinkConnector.
Ниже приведены ссылки на открытые репозитории некоторых коннекторов:
-
-
При необходимости создайте пользовательскую директорию для хранения плагинов (например, /var/lib/kafka-connect/jars ).
-
На каждом хосте с компонентом Kafka Connect Worker скопируйте в созданную директорию (или в существующую директорию /usr/lib/kafka-connect/plugins) плагин коннектора, например:
$ scp /tmp/clickhouse-kafka-connect-v1.0.16-confluent.jar olga@10.92.38.105:/var/lib/kafka-connect/jars -
Откройте страницу конфигурирования сервиса Kafka Connect, установите флаг Advanced, раскройте группу connect-distributed.properties→plugin.path, где уже указан путь к плагинам, добавленным по умолчанию — /usr/lib/kafka-connect/plugins.
Настройка plugin.path -
Если новый плагин сохранен в пользовательской директории, выберите поле Add property, и для новой строки значения параметра plugin.path укажите путь к пользовательской директории для хранения плагинов.
Указание нового plugin.path -
Перезагрузите сервис Kafka Connect при помощи действия Restart, нажав на иконку
в столбце Actions.
-
Проверьте, что новые настройки plugin.path отображаются в файле /etc/kafka-connect/config/connect-distributed.properties на хосте. Ниже приведен пример файла после внесенных изменений:
connect-distributed.properties# Maintained by ADCM # Kafka Broker Configuration bootstrap.servers=sov-ads-1.ru-central1.internal:9092 security.protocol=PLAINTEXT sasl.mechanism=none producer.sasl.mechanism=none consumer.sasl.mechanism=none consumer.security.protocol=PLAINTEXT producer.security.protocol=PLAINTEXT config.storage.topic=mm-connect-configs offset.storage.topic=mm-connect-offsets status.storage.topic=mm-connect-status rest.advertised.host.name = sov-ads-1.ru-central1.internal rest.advertised.listener=http listeners=http://0.0.0.0:8083 config.storage.replication.factor=1 connector.client.config.override.policy=None group.id=mm-connect key.converter=org.apache.kafka.connect.converters.ByteArrayConverter offset.flush.interval.ms=10000 offset.storage.replication.factor=1 plugin.path=/usr/lib/kafka-connect/plugins,/var/lib/kafka-connect/jars rest.port=8083 status.storage.replication.factor=1 value.converter=org.apache.kafka.connect.converters.ByteArrayConverter -
Проверьте, что при создании коннектора в пользовательском интерфейсе ADS Control отображается новый плагин.
Новый плагин коннектора в ADS ControlНаличие плагина также можно проверить при помощи HTTP-запроса к REST API любого из Kafka Connect Worker:
$ curl -X GET 'http://10.92.38.105:8083/connector-plugins' |jqВ ответ на запрос выводится список доступных плагинов для кластера и их версии.
Плагины коннекторов[ { "class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector", "type": "sink", "version": "v1.0.16" }, { "class": "org.apache.iceberg.connect.IcebergSinkConnector", "type": "sink", "version": "1.10.1.1-4.3.0-1" }, { "class": "org.apache.kafka.connect.tools.MockSinkConnector", "type": "sink", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.tools.VerifiableSinkConnector", "type": "sink", "version": "4.1.2.2-4.0.0-1" }, { "class": "io.debezium.connector.postgresql.PostgresConnector", "type": "source", "version": "3.5.1.1-4.0.0-1" }, { "class": "io.debezium.connector.sqlserver.SqlServerConnector", "type": "source", "version": "3.5.1.1-4.0.0-1" }, { "class": "org.apache.kafka.connect.mirror.MirrorCheckpointConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.mirror.MirrorHeartbeatConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.tools.MockSourceConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.tools.SchemaSourceConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" }, { "class": "org.apache.kafka.connect.tools.VerifiableSourceConnector", "type": "source", "version": "4.1.2.2-4.0.0-1" } ]
Влияние параметров Kafka Connect
Ниже описаны параметры Kafka Connect, влияющие на создание коннектора:
-
plugin.path — директории для размещения плагинов (по умолчанию
/usr/lib/kafka-connect/plugins); -
rest.port— порт REST API для формирования HTTP-запросов (по умолчанию8083); -
connector.client.config.override.policy — определяет, разрешены ли индивидуальные переопределения свойств клиента/безопасности для каждого коннектора (по умолчанию
None).
Указанные параметры записаны в файле /etc/kafka-connect/config/connect-distributed.properties на каждом хосте с компонентом Kafka Connect Worker.
Параметры могут быть изменены на странице конфигурирования сервиса Kafka Connect в группе connect-distributed.properties после установки флага Advanced.
После внесения всех изменений перезагрузите сервис Kafka Connect.
Ограничения при добавлении плагинов
-
После копирования JAR-файлов в указанную директорию требуется перезапуск сервиса Kafka Connect. Это связано с тем, что обнаружение плагинов происходит только при запуске воркера Kafka Connect.
-
При помощи инструментов ADS Control доступно только создание коннекторов в пользовательском интерфейсе. Добавление новых плагинов возможно только вручную на хосты с изменением параметров через ADCM.
-
В ADS 3.9.1.3 и ниже в кластере Kafka Connect может быть запущена только одна версия плагина.
-
В ADS 3.7.2.1 для создания Iceberg Sink Connector требуется установить вручную дополнительное значение plugin.path —
/usr/lib/kafka-connect/plugins. Использование механизма CLASSPATH не рекомендуется (конфликты библиотек, отсутствие изоляции).Начиная с ADS 3.9.0.1 значение
/usr/lib/kafka-connect/pluginsдля plugin.path предустановлено заранее, вместе с/var/lib/kafka-connect/libs.Начиная с ADS 4.0.0 значение
/usr/lib/kafka-connect/pluginsявляется единственным предустановленным значением plugin.path. По этому пути размещены все доступные по умолчанию в ADS плагины.