Добавление нового коннектора для Kafka Connect

Обзор

По умолчанию для создания в ADS доступны следующие коннекторы Kafka Connect:

Плагины этих коннекторов отображаются при создании коннектора в пользовательском интерфейсе ADS Control и размещаются в директории /usr/lib/kafka-connect/plugins, указанной в качестве значения параметра plugin.path в группе connect-distributed.properties на странице конфигурирования сервиса Kafka Connect.

Остальные коннекторы могут быть добавлены самостоятельно.

Начиная с ADS 4.0.0 (версия Kafka Connect 4.1.1 и выше) поддерживается запуск нескольких версий одного и того же плагина в одном кластере Kafka Connect. При создании коннектора доступно указание используемой версии плагина при помощи параметра connector.plugin.version. Доступные версии плагинов отображаются в ответ на запрос GET http://<worker-host>:8083/connector-plugins.

В ADS 3.9.1.3 и ниже в кластере Kafka Connect может быть запущена только одна версия плагина (выбирается последняя версия по порядку).

Добавление плагина коннектора

  1. Подготовьте один из вариантов плагина коннектора:

    • исполняемый JAR-файл (uber-JAR), содержащий Java-код коннектора и все его зависимости;

    • поддиректория плагина со всеми JAR-файлами и зависимостями (рекомендуется одна поддиректория на каждый плагин).

    ВНИМАНИЕ

    Плагин не должен включать runtime-библиотеки Kafka Connect.

    Для примера в данной статье используется исполняемый JAR-файл плагина для ClickHouseSinkConnector.

    Ниже приведены ссылки на открытые репозитории некоторых коннекторов:

  2. При необходимости создайте пользовательскую директорию для хранения плагинов (например, /var/lib/kafka-connect/jars ).

  3. На каждом хосте с компонентом 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
  4. Откройте страницу конфигурирования сервиса Kafka Connect, установите флаг Advanced, раскройте группу connect-distributed.properties→plugin.path, где уже указан путь к плагинам, добавленным по умолчанию — /usr/lib/kafka-connect/plugins.

    Настройка plugin.path
    Настройка plugin.path
  5. Если новый плагин сохранен в пользовательской директории, выберите поле Add property, и для новой строки значения параметра plugin.path укажите путь к пользовательской директории для хранения плагинов.

    Указание нового plugin.path
    Указание нового plugin.path
  6. Перезагрузите сервис Kafka Connect при помощи действия Restart, нажав на иконку actions default dark actions default light в столбце Actions.

  7. Проверьте, что новые настройки 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
  8. Проверьте, что при создании коннектора в пользовательском интерфейсе ADS Control отображается новый плагин.

    Новый плагин коннектора в 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 плагины.

Нашли ошибку? Выделите текст и нажмите Ctrl+Enter чтобы сообщить о ней