【问题标题】:Kafka Connect can't find connectorKafka Connect 找不到连接器
【发布时间】:2019-04-24 01:19:40
【问题描述】:

我正在尝试使用 Kafka Connect Elasticsearch 连接器,但没有成功。它因以下错误而崩溃:

[2018-11-21 14:48:29,096] ERROR Stopping after connector error (org.apache.kafka.connect.cli.ConnectStandalone:108)
java.util.concurrent.ExecutionException: org.apache.kafka.connect.errors.ConnectException: Failed to find any class that implements Connector and which name matches io.confluent.connect.elasticsearch.ElasticsearchSinkConnector , available connectors are: PluginDesc{klass=class org.apache.kafka.connect.file.FileStreamSinkConnector, name='org.apache.kafka.connect.file.FileStreamSinkConnector', version='1.0.1', encodedVersion=1.0.1, type=sink, typeName='sink', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.file.FileStreamSourceConnector, name='org.apache.kafka.connect.file.FileStreamSourceConnector', version='1.0.1', encodedVersion=1.0.1, type=source, typeName='source', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.MockConnector, name='org.apache.kafka.connect.tools.MockConnector', version='1.0.1', encodedVersion=1.0.1, type=connector, typeName='connector', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.MockSinkConnector, name='org.apache.kafka.connect.tools.MockSinkConnector', version='1.0.1', encodedVersion=1.0.1, type=sink, typeName='sink', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.MockSourceConnector, name='org.apache.kafka.connect.tools.MockSourceConnector', version='1.0.1', encodedVersion=1.0.1, type=source, typeName='source', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.SchemaSourceConnector, name='org.apache.kafka.connect.tools.SchemaSourceConnector', version='1.0.1', encodedVersion=1.0.1, type=source, typeName='source', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.VerifiableSinkConnector, name='org.apache.kafka.connect.tools.VerifiableSinkConnector', version='1.0.1', encodedVersion=1.0.1, type=source, typeName='source', location='classpath'}, PluginDesc{klass=class org.apache.kafka.connect.tools.VerifiableSourceConnector, name='org.apache.kafka.connect.tools.VerifiableSourceConnector', version='1.0.1', encodedVersion=1.0.1, type=source, typeName='source', location='classpath'}

我在 kafka 子文件夹中解压缩了插件的构建,并且在 connect-standalone.properties 中有以下行:

plugin.path=/opt/kafka/plugins/kafka-connect-elasticsearch-5.0.1/src/main/java/io/confluent/connect/elasticsearch

我可以看到该文件夹​​中的各种连接器,但 Kafka Connect 没有加载它们;但它确实会加载标准连接器,如下所示:

[2018-11-21 14:56:28,258] INFO Added plugin 'org.apache.kafka.connect.transforms.Cast$Value' (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:136)
[2018-11-21 14:56:28,259] INFO Added aliases 'FileStreamSinkConnector' and 'FileStreamSink' to plugin 'org.apache.kafka.connect.file.FileStreamSinkConnector' (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:335)
[2018-11-21 14:56:28,260] INFO Added aliases 'FileStreamSourceConnector' and 'FileStreamSource' to plugin 'org.apache.kafka.connect.file.FileStreamSourceConnector' (org.apache.kafka.connect.runtime.isolation.DelegatingClassLoader:335)

如何正确注册连接器?

【问题讨论】:

    标签: apache-kafka apache-kafka-connect confluent-platform


    【解决方案1】:

    我昨天在 docker 中的 kafka 上手动运行了 jdbc 连接器,没有融合平台等,只是为了了解这些东西在下面是如何工作的。我不必在我身边建造罐子或任何类似的东西。希望它与您相关 - 我所做的是(我将跳过 docker 部分如何使用连接器安装目录等):

    • https://www.confluent.io/connector/kafka-connect-jdbc/下载连接器,解压压缩包
    • 将 zip 的内容放入属性文件中配置的路径中的目录中(如下图第 3 点所示)-

      plugin.path=/plugins
      

      所以树看起来像这样:

      /plugins/
      └── jdbcconnector
          └──assets
          └──doc
          └──etc
          └──lib
      

      注意 lib 目录的依赖关系,其中之一是 kafka-connect-jdbc-5.0.0.jar

    • 现在您可以尝试运行连接器

      ./connect-standalone.sh connect-standalone.properties jdbc-connector-config.properties
      

      connect-standalone.properties 是 kafka-connect 所需的常用属性,在我的例子中:

      bootstrap.servers=localhost:9092
      key.converter=org.apache.kafka.connect.json.JsonConverter
      value.converter=org.apache.kafka.connect.json.JsonConverter
      key.converter.schemas.enable=true
      value.converter.schemas.enable=true
      offset.storage.file.filename=/tmp/connect.offsets
      offset.flush.interval.ms=10000
      plugin.path=/plugins
      rest.port=8086
      rest.host.name=127.0.0.1
      

      jdbc-connector-config.properties 涉及更多,因为它只是此特定连接器的配置,您需要深入研究连接器文档 - 对于 jdbc 源它是 https://docs.confluent.io/current/connect/kafka-connect-jdbc/source-connector/source_config_options.html

    【讨论】:

    • 对于 JDBC Connect,我注意到驱动程序甚至可以放在它的子文件夹中。例如 lib/drivers
    • 当我尝试执行 plugin.path 时,我不断收到“没有这样的文件或目录”,知道为什么吗?我对此完全陌生
    【解决方案2】:

    编译后的 JAR 需要可供 Kafka Connect 使用。你有几个选择:

    1. 使用 Confluent 平台,其中包括预构建的 Elasticsearch(和其他):https://www.confluent.io/download/。有 zip、rpm/deb、Docker 镜像等可用。

    2. 自己构建 JAR。这通常涉及:

      cd kafka-connect-elasticsearch-5.0.1
      mvn clean package
      

      然后将生成的 kafka-connect-elasticsearch-5.0.1.jar JAR 放入 Kafka Connect 中配置的路径中,并使用 plugin.path

    您可以在此处找到有关使用 Kafka Connect 的更多信息:

    免责声明:我为 Confluent 工作,并撰写了上述博客文章。

    【讨论】:

    • 自己构建会导致:[ERROR] [ERROR] Some problems were encountered while processing the POMs: [FATAL] Non-resolvable parent POM for io.confluent:kafka-connect-elasticsearch:[unknown-version]: Could not transfer artifact io.confluent:common:pom:5.0.1 from/to confluent (${confluent.maven.repo}):
    • ¯_(ツ)_/¯ 这是在 Confluent Platform 中使用预构建版本的一个很好的理由 ;) 它是开源的,可以免费使用。如果您真的不想这样做,您可以 d/l 它并简单地提取 JAR 并将其部署到您现有的安装中。
    • 好的,我已将预构建版本保存在 /var/confluentinc-kafka-connect-elasticsearch-5.0.0/ 中。在我的配置中,我有这一行:plugin.path=/var/silverbolt/confluentinc-kafka-connect-elasticsearch-5.0.0/ 。我仍然收到关于没有匹配的连接器类的相同错误。
    【解决方案3】:

    插件路径必须加载包含编译代码的 JAR 文件,而不是源代码的原始 Java 类 (src/main/java)。

    它还需要是包含这些插件的其他目录的父目录

    plugin.path=/opt/kafka-connect/plugins/
    

    在哪里

    $ ls - lR /opt/kafka-connect/plugins/
    kafka-connect-elasticsearch-x.y.z/
        file1.jar
        file2.jar 
        etc
    

    参考 - Manually installing Community Connectors

    Confluent 平台中的 Kafka Connect 启动脚本也会自动(习惯于?)读取所有匹配 share/java/kafka-connect-* 的文件夹,所以这是一种方法。至少,如果你在插件路径中包含 Confluent 包安装的share/java 文件夹的路径,它将继续这样做

    如果你对 Maven 不是很熟悉,或者即使你很熟悉,那么你实际上不能只克隆第一个 Kafka 的 Elasticsearch 连接器 repo 和 build the master branch; it has prerequisites,然后首先克隆常见的 Confluent repo。否则,您必须签出与 Confluent 版本匹配的 Git 标记,例如 5.0.1-post

    一个更简单的选择是grab the package using Confluent Hub CLI

    如果这些都不起作用,那么只需下载 Confluent 平台并使用 Kafka Connect 脚本将是最简单的。这并不意味着您需要使用其中的 Kafka 或 Zookeeper 配置

    【讨论】:

    • 我成功构建了 JAR,将其移动到 /plugins/ 下的 a 文件夹中,将路径添加到配置中,但仍然出现相同的“找不到任何类”错误。
    • 我相信 Elastic 连接器实际上会生成一个您需要提取的 tar.gz 文件。它不会只创建一个包含所有所需类的 JAR
    • 我正在查看target 文件夹,它创建了一个kafka-connect-elasticsearch-5.0.1.jar 文件。没有tar.gz,我可以看到。
    • 对我来说,在末尾添加逗号解决了plugin.path=/opt/bitnami/kafka/connectors, 的问题。没有逗号,kafka一直抱怨找不到类。
    猜你喜欢
    • 2022-05-31
    • 2018-08-27
    • 2021-01-13
    • 2019-06-17
    • 2020-08-05
    • 2021-08-28
    • 2020-06-19
    • 2018-02-06
    • 2020-08-26
    相关资源
    最近更新 更多