【问题标题】:How to stream data from Kafka to MongoDB by Kafka Connector如何通过 Kafka 连接器将数据从 Kafka 流式传输到 MongoDB
【发布时间】:2019-11-14 18:11:36
【问题描述】:

我想使用 Kafka 连接器将数据从 Kafka 流式传输到 MongoDB。 我找到了这个https://github.com/hpgrahsl/kafka-connect-mongodb。但是没有步骤可做。

谷歌搜索后,似乎导致我不想使用的 Confluent Platform。

谁能分享文档/指南如何在不使用 Confluent 平台或其他 Kafka 连接器将数据从 Kafka 流式传输到 MongoDB 的情况下使用 kafka-connect-mongodb

提前谢谢你。


我尝试了什么

Step1:我从maven central下载mongo-kafka-connect-0.1-all.jar

Step2:将jar文件复制到kafka里面的一个新文件夹plugins(我在Windows上使用Kafka,所以目录是D:\git\1.libraries\kafka_2.12-2.2.0\plugins

第三步:编辑文件connect-standalone.properties,添加新行 plugin.path=/git/1.libraries/kafka_2.12-2.2.0/plugins

Step4:我为 mongoDB sink MongoSinkConnector.properties 添加新的配置文件

name=mongo-sink
topics=test
connector.class=com.mongodb.kafka.connect.MongoSinkConnector
tasks.max=1
key.ignore=true

# Specific global MongoDB Sink Connector configuration
connection.uri=mongodb://localhost:27017,mongo1:27017,mongo2:27017,mongo3:27017
database=test_kafka
collection=transaction
max.num.retries=3
retries.defer.timeout=5000
type.name=kafka-connect

Step5:运行命令bin\windows\connect-standalone.bat config\connect-standalone.properties config\MongoSinkConnector.properties

但是,我得到了错误

[2019-07-09 10:19:09,466] WARN The configuration 'offset.flush.interval.ms' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,467] WARN The configuration 'key.converter.schemas.enable' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,467] WARN The configuration 'offset.storage.file.filename' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,468] WARN The configuration 'value.converter.schemas.enable' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,469] WARN The configuration 'plugin.path' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,469] WARN The configuration 'value.converter' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
[2019-07-09 10:19:09,470] WARN The configuration 'key.converter' was supplied but isn't a known config. (org.apache.kafka.clients.admin.AdminClientConfig)
Jul 09, 2019 10:19:10 AM org.glassfish.jersey.internal.inject.Providers checkProviderRuntime
WARNING: A provider org.apache.kafka.connect.runtime.rest.resources.ConnectorPluginsResource registered in SERVER runtime does not implement any provider interfaces applicable in the SERVER runtime. Due to constraint configuration problems the provider org.apache.kafka.connect.runtime.rest.resources.ConnectorPluginsResource will be ignored.
Jul 09, 2019 10:19:10 AM org.glassfish.jersey.internal.inject.Providers checkProviderRuntime
WARNING: A provider org.apache.kafka.connect.runtime.rest.resources.RootResource registered in SERVER runtime does not implement any provider interfaces applicable in the SERVER runtime. Due to constraint configuration problems the provider org.apache.kafka.connect.runtime.rest.resources.RootResource will be ignored.
Jul 09, 2019 10:19:10 AM org.glassfish.jersey.internal.inject.Providers checkProviderRuntime
WARNING: A provider org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource registered in SERVER runtime does not implement any provider interfaces applicable in the SERVER runtime. Due to constraint configuration problems the provider org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource will be ignored.
Jul 09, 2019 10:19:11 AM org.glassfish.jersey.internal.Errors logErrors
WARNING: The following warnings have been detected: WARNING: The (sub)resource method listConnectors in org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource contains empty path annotation.
WARNING: The (sub)resource method createConnector in org.apache.kafka.connect.runtime.rest.resources.ConnectorsResource contains empty path annotation.
WARNING: The (sub)resource method listConnectorPlugins in org.apache.kafka.connect.runtime.rest.resources.ConnectorPluginsResource contains empty path annotation.
WARNING: The (sub)resource method serverInfo in org.apache.kafka.connect.runtime.rest.resources.RootResource contains empty path annotation.

[2019-07-09 10:19:12,302] ERROR WorkerSinkTask{id=mongo-sink-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask)
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:178)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:104)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:487)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:464)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:320)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error:
        at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:344)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$1(WorkerSinkTask.java:487)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:128)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:162)
        ... 13 more
Caused by: org.apache.kafka.common.errors.SerializationException: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'this': was expecting 'null', 'true', 'false' or NaN
 at [Source: (byte[])"this is a message"; line: 1, column: 6]
Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'this': was expecting 'null', 'true', 'false' or NaN
 at [Source: (byte[])"this is a message"; line: 1, column: 6]
        at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:1804)
        at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:703)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3532)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3508)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._matchToken2(UTF8StreamJsonParser.java:2843)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._matchTrue(UTF8StreamJsonParser.java:2777)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:807)
        at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:729)
        at com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4042)
        at com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:2571)
        at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:50)
        at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:342)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$1(WorkerSinkTask.java:487)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:128)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:162)
        at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:104)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:487)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:464)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:320)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
        at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
        at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
[2019-07-09 10:19:12,305] ERROR WorkerSinkTask{id=mongo-sink-0} Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask)

我设置错了什么配置或者我错过了什么?


我修好了。现在,我可以成功地将数据从 Kafka 流式传输到 MongoDB

我的解决方法是:

  1. 将我的卡夫卡移动到C:\kafka_2.12-2.2.0
  2. 更新plugin_path对应新路径
  3. 更新配置文件connect-standalone.properties

【问题讨论】:

  • 出于兴趣为什么不想使用 Confluent 平台?
  • @RobinMoffatt 我担心它的商业功能“单个 Kafka 代理永久免费。无限 Kafka 代理的 30 天免费试用”。这是否意味着我不能创建超过 1 个具有免费功能的容错代理?
  • 注意:Confluent 平台的所有组件都可以与现有的 Kafka 安装一起使用。 Kafka Connect 不是特定于 Confluent 平台的
  • @VuLeAnh 这只是企业许可证功能 - Confluent 控制中心、复制器等。 Confluent 平台的许多组件都获得了社区许可。
  • 查看confluent.io/download查看对比图表。

标签: mongodb apache-kafka apache-kafka-connect


【解决方案1】:

MongoDB 本身有一个官方的源和接收器连接器。它在 Confluent Hub 上可用:https://www.confluent.io/hub/mongodb/kafka-connect-mongodb

如果您不想使用 Confluent Platform,您可以自己部署 Apache Kafka - 它已经包含 Kafka Connect。您使用哪些插件(连接器)取决于您。在这种情况下,您将使用 Kafka Connect(Apache Kafka 的一部分)和 kafka-connect-mongodb(由 MongoDB 提供)。

如何使用它的文档在这里:https://docs.mongodb.com/kafka-connector/current/

【讨论】:

    【解决方案2】:

    尽管这个问题有点老了。这是我在 ubuntu 系统上将 kafka_2.12-2.6.0 连接到 mongodb(4.4 版)的方法:

    一个。从 here 下载 mongodb 连接器 '*-all.jar' 。最后带有 'all' 的 Mongodb-kafka 连接器也将包含所有连接器依赖项。

    湾。将此 jar 文件放入您的 kafka 的 lib 文件夹中

    C。将“connect-standalone_bare.properties”配置为:

    bootstrap.servers=localhost:9092
    key.converter=org.apache.kafka.connect.json.JsonConverter
    value.converter=org.apache.kafka.connect.json.JsonConverter
    key.converter.schemas.enable=false
    value.converter.schemas.enable=false
    offset.storage.file.filename=/tmp/connect.offsets
    offset.flush.interval.ms=10000
    

    d。将“MongoSinkConnector.properties”配置为:

    name=mongo-sink
    topics=test
    connector.class=com.mongodb.kafka.connect.MongoSinkConnector
    tasks.max=1
    key.ignore=true
    connection.uri=mongodb://localhost:27017
    database=test_kafka
    collection=transaction
    max.num.retries=3
    retries.defer.timeout=5000
    type.name=kafka-connect
    schemas.enable=false
    

    将两个“properties”文件放在此处:$HOME/Documents/kafka/config

    e。启动connector-process,如下:

    export folder_path="$HOME/Documents/kafka/config"
    connect-standalone.sh  $folder_path/connect-standalone_bare.properties $folder_path/MongoSinkConnector.properties
    

    e。在 kafka 中,启动 zookeeper-server 和 kafka-server。创建主题“测试”。在 mongod 服务器中,创建数据库“test_kafka”并在其下创建一个集合“transaction”。

    F。启动 kafka 生产者:

    kafka-console-producer.sh --broker-list localhost:9092  --topic test
    

    并输入:{"abc" : "def" }

    您应该可以在 mongodb (db.transaction.find() ) 中看到它。

    【讨论】:

    • mongodb作为源怎么办?
    • 在e部分,不应该把语句改成start zookeeper-server和kafka-server吗?
    • 这不起作用。无论生产者写什么,我都无法在数据库中看到。
    • connect-standalone_bare.properties 在哪里?我正在使用泊坞窗
    猜你喜欢
    • 2017-08-22
    • 2023-03-26
    • 1970-01-01
    • 2021-11-26
    • 2021-12-23
    • 2020-04-27
    • 1970-01-01
    • 2017-09-15
    • 1970-01-01
    相关资源
    最近更新 更多