如何将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的要求
原配置问题分析:
- 直接设置
ExtractField.field=id时,提取的是包含$oid的Map结构,而Elasticsearch Sink不支持将Map作为document_id,因此触发MAP is not supported as the document id错误 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ê
相关产品推荐
相关产品推荐

