如何用Kafka Connect Elasticsearch Sink转换嵌套JSON字段并生成随机ID?
解决方案
你遇到的问题核心有两点:一是配置里的字段名拼写错误(arrtubites应为attributes),二是ReplaceField转换仅支持顶层字段的重命名,无法直接处理嵌套字段路径。结合你的需求,以下是两种可行的配置方案:
方案一:使用Jolt Transform(推荐,灵活处理嵌套结构)
Jolt是专门用于JSON结构转换的工具,能轻松实现嵌套字段的映射与扁平化。Confluent Platform默认包含该插件,开源Kafka Connect需单独安装Jolt Transform插件。
完整配置如下:
{ "name": "elasticsearch-sink", "config": { "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "tasks.max": "1", "topics": "document, document-processing-error", "key.ignore": "true", "schema.ignore": "true", "connection.url": "http://elasticsearch:9200", "type.name": "_doc", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "flush.synchronously": "true", "transforms": "joltTransform, appendTimestampToIX", // Jolt转换:将attributes.3映射为顶层cashier字段,自动丢弃原嵌套结构 "transforms.joltTransform.type": "org.apache.kafka.connect.transforms.JoltTransform$Value", "transforms.joltTransform.spec": "[{\"operation\": \"shift\", \"spec\": {\"attributes\": {\"3\": \"cashier\"}}}]"", // 保留原有的时间戳路由配置 "transforms.appendTimestampToIX.type": "org.apache.kafka.connect.transforms.TimestampRouter", "transforms.appendTimestampToIX.topic.format": "${topic}-${timestamp}", "transforms.appendTimestampToIX.timestamp.format": "yyyy-MM-dd" } }
关键说明:
- Jolt的
shift操作直接将嵌套的attributes.3值映射到顶层cashier字段,无需额外处理原结构 - 由于配置了
key.ignore: true,当Kafka消息无key时,Elasticsearch会自动生成随机文档ID,完全匹配你的需求
方案二:使用内置Transform组合(无需额外插件)
如果无法使用Jolt插件,可通过ExtractField+WrapField+ReplaceField的组合实现字段转换:
{ "name": "elasticsearch-sink", "config": { "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "tasks.max": "1", "topics": "document, document-processing-error", "key.ignore": "true", "schema.ignore": "true", "connection.url": "http://elasticsearch:9200", "type.name": "_doc", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "flush.synchronously": "true", "transforms": "extractAttr3, wrapAsTempField, renameToCashier, appendTimestampToIX", // 提取attributes.3字段的值,此时消息value变为字符串"Mike" "transforms.extractAttr3.type": "org.apache.kafka.connect.transforms.ExtractField$Value", "transforms.extractAttr3.field": "attributes.3", // 将字符串值包装为临时对象{"3": "Mike"} "transforms.wrapAsTempField.type": "org.apache.kafka.connect.transforms.WrapField$Value", "transforms.wrapAsTempField.field": "3", // 将临时字段"3"重命名为cashier,得到{"cashier": "Mike"} "transforms.renameToCashier.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.renameToCashier.renames": "3:cashier", // 保留原有的时间戳路由配置 "transforms.appendTimestampToIX.type": "org.apache.kafka.connect.transforms.TimestampRouter", "transforms.appendTimestampToIX.topic.format": "${topic}-${timestamp}", "transforms.appendTimestampToIX.timestamp.format": "yyyy-MM-dd" } }
额外注意:
你的Elasticsearch索引模板中index_patterns为documents-*,但生成的索引是document-2022-11-03,两者不匹配,需修改模板的index_patterns为document-*,否则索引将不会应用该模板的映射规则。
内容的提问来源于stack exchange,提问作者ArmanDinagur
相关产品推荐
相关产品推荐

