You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

MongoDB Kafka Source Connector提取文档ID作为消息Key失败求助

问题解决:MongoDB源连接器提取文档ID作为Kafka消息Key

核心错误原因

你当前的ExtractField转换配置直接尝试提取$oid字段,但原始Key的顶层结构中不存在该字段;同时MongoDB的ObjectId在Kafka Connect中是特殊类型,并非普通嵌套结构体,无法直接通过ExtractField提取内部的oid字符串,这才触发了Schema type of struct but value was of type: object_id异常。

修正后的连接器配置

{
  "name": "new-items-source",
  "config": {
    "name": "new-items-source",
    "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector",
    "tasks.max": "1",
    "connection.uri": "mongodb://mongo1:27017",
    "database": "store",
    "collection": "items",
    "topic.prefix": "new_items",
    "publish.full.document.only": "true",
    "pipeline": "[{\"$match\": {\"operationType\": \"insert\"}}, {\"$project\": {\"fullDocument.bar\": 0, \"fullDocument._class\": 0, \"fullDocument.foo\": 0, \"fullDocument._id\": 0 } } ]",
    "transforms": "extractDocKey,extractId,castToString",
    "transforms.extractDocKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
    "transforms.extractDocKey.field": "documentKey",
    "transforms.extractId.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
    "transforms.extractId.field": "_id",
    "transforms.castToString.type": "org.apache.kafka.connect.transforms.Cast$Key",
    "transforms.castToString.spec": "string",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter.schema.registry.url": "http://schemaregistry0:8585",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schemas.enable": "true",
    "output.schema.value": "{\"type\": \"record\", \"name\": \"NewItemMessage\", \"namespace\": \"com.acme.store.adapters.out.kafka.avro\", \"fields\": [{\"name\": \"author\", \"type\": \"string\"}, {\"name\": \"state\", \"type\": \"string\"}, {\"name\": \"text\", \"type\": \"string\"}] }",
    "output.format.key": "string",
    "output.format.value": "schema",
    "output.json.formatter": "com.mongodb.kafka.connect.source.json.formatter.ExtendedJson",
    "errors.log.include.messages": "true",
    "errors.log.enable": "true"
  }
}

关键修改说明

  1. Pipeline优化
    在$project阶段添加"fullDocument._id": 0,直接从源数据中移除不需要的_id字段,确保最终消息体只保留你需要的author、text、state字段。

  2. 多步骤Key转换

    • extractDocKey:从原始事件Key中提取documentKey字段(该字段存储了文档的唯一ID信息)
    • extractId:从documentKey中提取_id字段(此时为MongoDB ObjectId类型)
    • castToString:将ObjectId类型强制转换为字符串,得到纯ID字符串作为Kafka消息的Key
  3. Converter调整
    因为最终Key是字符串类型,使用StringConverter更简洁,无需依赖Schema Registry;如果坚持使用Avro Converter,需将output.schema.key设置为{"type": "string"}并保留Key的Schema Registry配置。

  4. Value Schema更新
    修改output.schema.value,移除原有的_id字段,匹配你期望的消息体结构。

内容的提问来源于stack exchange,提问作者Rubén

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 04:15:01