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" } }
关键修改说明
Pipeline优化
在$project阶段添加"fullDocument._id": 0,直接从源数据中移除不需要的_id字段,确保最终消息体只保留你需要的author、text、state字段。多步骤Key转换
extractDocKey:从原始事件Key中提取documentKey字段(该字段存储了文档的唯一ID信息)extractId:从documentKey中提取_id字段(此时为MongoDB ObjectId类型)castToString:将ObjectId类型强制转换为字符串,得到纯ID字符串作为Kafka消息的Key
Converter调整
因为最终Key是字符串类型,使用StringConverter更简洁,无需依赖Schema Registry;如果坚持使用Avro Converter,需将output.schema.key设置为{"type": "string"}并保留Key的Schema Registry配置。Value Schema更新
修改output.schema.value,移除原有的_id字段,匹配你期望的消息体结构。
内容的提问来源于stack exchange,提问作者Rubén
相关产品推荐
相关产品推荐

