【问题标题】:Error handling for invalid JSON in kafka sink connectorkafka sink 连接器中无效 JSON 的错误处理
【发布时间】:2020-05-26 21:05:29
【问题描述】:

我有一个用于 mongodb 的接收器连接器,它从主题中获取 json 并将其放入 mongoDB 集合中。但是,当我从生产者向该主题发送无效的 JSON(例如,使用无效的特殊字符“)=> {"id":1,"name":"\"} 时,连接器停止。我尝试使用errors.tolerance = all,但同样的事情正在发生。什么应该发生的是连接器应该跳过并记录无效的 JSON,并保持连接器运行。我的分布式模式连接器如下:

{
  "name": "sink-mongonew_test1",  
  "config": {
    "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector",
    "topics": "error7",
    "connection.uri": "mongodb://****:27017", 
    "database": "abcd", 
    "collection": "abc",
    "type.name": "kafka-connect",
    "key.ignore": "true",

    "document.id.strategy": "com.mongodb.kafka.connect.sink.processor.id.strategy.PartialValueStrategy",
    "value.projection.list": "id",
    "value.projection.type": "whitelist",
    "writemodel.strategy": "com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneTimestampsStrategy",

    "delete.on.null.values": "false",

    "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",

    "errors.tolerance": "all",
    "errors.log.enable": "true",
    "errors.log.include.messages": "true",
    "errors.deadletterqueue.topic.name": "crm_data_deadletterqueue",
    "errors.deadletterqueue.topic.replication.factor": "1",
    "errors.deadletterqueue.context.headers.enable": "true"
  }
}

【问题讨论】:

  • 你为什么要生成无效的 json?如何?如果你使用任何 json 库,它不会产生无效的 json 字符串
  • 您使用的是哪个版本的 Connect?
  • @cricket_007 我正在构建一个 kafka-streams 应用程序,其中存在无效 JSON 的可能性。 connect是connect-api-1.0.1.3.0.0.0-1634.jar版本,kafka是3.0
  • 再次,不清楚您的生产者是如何创建无效 JSON 的。 Kafka Streams 有自己的错误处理,无论如何您都应该过滤掉无效记录……Kafka 没有 3.0 版。相对而言,Confluent Platform 3.0 真的很老了。 HDP 3.0 使用的是 Kafka 1.x,我认为……如果你安装了 HDF,Cloudera 会建议你使用 Nifi
  • 你读过这个吗? confluent.io/blog/…

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


【解决方案1】:

Apache Kafka 2.0 以来,Kafka Connect 已包含错误处理选项,包括将消息路由到死信队列的功能,这是构建数据管道的常用技术。

https://www.confluent.io/blog/kafka-connect-deep-dive-error-handling-dead-letter-queues/

正如评论所言,您使用的是 connect-api-1.0.1.*.jar,版本 1.0.1,这就解释了为什么这些属性不起作用

除了运行较新版本的 Kafka Connect 之外,您还可以选择 Nifi 或 Spark Structured Streaming

【讨论】:

  • 这是我应该在 kafka 的 libs 文件夹中替换(使用最新版本)的唯一 jar 吗?
  • 没有。导入诸如 kafka-clients 和 kafka_2.11 和 zookeeper-3.4.13、connect-utils 等 JAR 文件...
  • 您不需要使用 Hortonworks 提供的 Kafka Connect。只需在同一台机器上下载最新的 Kafka 库并为 Connect Distributed 配置更新的属性文件,然后针对现有的 Kafka 引导服务器运行它。 Cloudera 无论如何都没有正式支持 Kafka Connect API,所以你真的是靠自己
猜你喜欢
  • 2022-09-24
  • 2020-07-26
  • 1970-01-01
  • 2019-06-22
  • 2021-05-05
  • 2019-12-20
  • 2019-06-15
  • 2019-09-18
  • 2020-11-04
相关资源
最近更新 更多