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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 14:12:05