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

如何将Kafka Key中的嵌套JSON值作为Elasticsearch的document_id

解决方案

要提取嵌套的_id.$oid作为Elasticsearch的document_id,你需要先通过Flatten转换将Key中的嵌套结构扁平化,再提取目标字段。以下是修正后的Sink配置:

transforms=FlattenKey,ExtractOid
transforms.FlattenKey.type=org.apache.kafka.connect.transforms.Flatten$Key
transforms.FlattenKey.delimiter=.
transforms.ExtractOid.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.ExtractOid.field=_id.$oid

配置说明:

  • FlattenKey:将Key中的嵌套JSON结构扁平化,原Key {"_id": {"$oid": "64e5cf730ed5822aa85b91ad"}} 会被转换为 {"_id.$oid": "64e5cf730ed5822aa85b91ad"}
  • ExtractOid:从扁平化后的Key中提取_id.$oid字段的值,得到纯字符串格式的ID,符合Elasticsearch对document_id的要求

原配置问题分析:

  1. 直接设置ExtractField.field=id时,提取的是包含$oid的Map结构,而Elasticsearch Sink不支持将Map作为document_id,因此触发MAP is not supported as the document id错误
  2. ExtractField转换不支持嵌套字段路径(如id.$oid),所以直接指定该路径会导致字段无法识别,触发Key is used as document id and can not be null错误

如果需要更简洁的字段名,可在Flatten后添加ReplaceField转换重命名字段:

transforms=FlattenKey,RenameOid,ExtractOid
transforms.FlattenKey.type=org.apache.kafka.connect.transforms.Flatten$Key
transforms.FlattenKey.delimiter=.
transforms.RenameOid.type=org.apache.kafka.connect.transforms.ReplaceField$Key
transforms.RenameOid.renames=_id.$oid:document_id
transforms.ExtractOid.type=org.apache.kafka.connect.transforms.ExtractField$Key
transforms.ExtractOid.field=document_id

内容的提问来源于stack exchange,提问作者Thanh Long Lê

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 21:35:13