【问题标题】:How to tell MongoSource (using Kafka Connect) what Key to serialize如何告诉 MongoSource(使用 Kafka Connect)序列化的密钥
【发布时间】:2020-04-14 22:44:43
【问题描述】:

我正在使用 mongo 源来监听 mongo 更改流并将所有事件放入 kafka,但我正在绞尽脑汁地寻找一种从事件中提取“Real”键的方法。我尝试了转换,但它不起作用,给我错误:

Caused by: org.apache.kafka.connect.errors.DataException: Only Struct objects supported for [copying fields from value to key], found: java.lang.String

在 Mongo 源中我发现了这个 line

这基本上意味着它甚至没有一些密钥处理,而是寻找“_id”字段(这不是文档的 id,它是一个恢复令牌信息)

相反,我想将主题的键设置为“documentKey”。

以下是连接器获取的事件示例:

{
 "_id": {
    "_data": "DSAD45543FFWEHTEY004....."
  },
  "operationType": "replace",
  "clusterTime": {
    "$timestamp": {
      "t": 1446707990,
      "i": 1
    }
  },
  "fullDocument": {
    "_id": {
      "$binary": "FxVFgHFRhrr/z+zUc/w==",
      "$type": "03"
    },
    ...
  },
  "ns": {
    "db": "somedb",
    "coll": "somecol"
  },
  "documentKey": {
    "_id": {
      "$binary": "FxVFgHFRhrr/z+zUc/w==",
      "$type": "03"
    }
  }
}

我使用了以下配置:

"transforms":"createKey",
"transforms.createKey.type":"org.apache.kafka.connect.transforms.ValueToKey",
"transforms.createKey.fields":"documentKey"

我试过了:

org.apache.kafka.connect.json.JsonConverter

还有StringConverter(虽然我不认为这可以用字符串来完成)

org.apache.kafka.connect.storage.StringConverter

有没有办法提取密钥? 请注意:架构已禁用。

【问题讨论】:

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


    【解决方案1】:

    这是因为 MongoDB Source Connector for Kafka 还不支持它。从 1.3 版开始,它应该支持高级密钥选择。

    https://jira.mongodb.org/browse/KAFKA-40

    【讨论】:

      【解决方案2】:

      请注意:架构已禁用

      在这种情况下,您不能使用 ValueToKey 转换。但是,即使您可以,该转换也不支持有效负载中的嵌套值,在您的情况下类似于documentKey._id.$binary

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2014-09-14
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-07-01
        • 2020-11-20
        • 1970-01-01
        • 2020-05-26
        相关资源
        最近更新 更多