Bor*_*s K 4 apache-kafka confluent apache-kafka-connect
我正在尝试使用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'}
Run Code Online (Sandbox Code Playgroud)
我已经在kafka子文件夹中解压缩了插件的构建,并且在connect-standalone.properties中包含以下行:
plugin.path=/opt/kafka/plugins/kafka-connect-elasticsearch-5.0.1/src/main/java/io/confluent/connect/elasticsearch
Run Code Online (Sandbox Code Playgroud)
我可以看到该文件夹中的各种连接器,但是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)
Run Code Online (Sandbox Code Playgroud)
如何正确注册连接器?
我昨天在没有融合平台等的docker上的kafka上手动运行了jdbc连接器,只是为了了解这些东西在下面如何工作。我不必在我这边或类似的东西上盖罐子。希望这对您有意义-我所做的是(我将跳过docker部件如何使用连接器等方式挂载dir):
将zip的内容放在属性文件中配置的路径中的目录中(如下所示,在第3点中)-
plugin.path=/plugins
Run Code Online (Sandbox Code Playgroud)
所以树看起来像这样:
/plugins/
??? jdbcconnector
???assets
???doc
???etc
???lib
Run Code Online (Sandbox Code Playgroud)
注意lib dir的依赖关系,其中之一是kafka-connect-jdbc-5.0.0.jar
现在您可以尝试运行连接器
./connect-standalone.sh connect-standalone.properties jdbc-connector-config.properties
Run Code Online (Sandbox Code Playgroud)
在我的情况下,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
Run Code Online (Sandbox Code Playgroud)
jdbc-connector-config.properties涉及更多,因为它只是为此特定连接器的配置,您需要深入研究连接器文档-对于jdbc源,它是https://docs.confluent.io/current/connect/kafka-connect -jdbc / source-connector / source_config_options.html
编译后的 JAR 需要可供 Kafka Connect 使用。您在这里有几个选择:
使用 Confluence Platform,其中包括预构建的 Elasticsearch(和其他):https: //www.confluence.io/download/。有 zip、rpm/deb、Docker 镜像等可用。
自己构建 JAR。这通常涉及:
cd kafka-connect-elasticsearch-5.0.1
mvn clean package
Run Code Online (Sandbox Code Playgroud)
然后获取生成的kafka-connect-elasticsearch-5.0.1.jarJAR 并将其放入 Kafka Connect with 中配置的路径中plugin.path。
您可以在此处找到有关使用 Kafka Connect 的更多信息:
免责声明:我在 Confluence 工作,并撰写了上述博客文章。
| 归档时间: |
|
| 查看次数: |
2922 次 |
| 最近记录: |