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

如何使用NiFi将Kafka偏移量、分区信息添加到消息元数据中

在NiFi中为Kafka消息添加offset和partition至metadata字段

原始消息结构

从Kafka读取的原始JSON消息如下:

{
  "body": {
    "metadata": {
      "id": "bce16e11"
    },
    "eventDetails": {
      "eventID": "c5f615f1",
      "customerId": "123456789",
      "Name": "NEW"
    }
  }
}

目标结构

需要将NiFi属性中的kafka.offset和kafka.partition添加到body.metadata字段中,最终结构如下:

{
  "body" : {
    "metadata" : {
      "id" : "bce16e11",
      "kafkaOffset" : 4537732,
      "kafkaPartition" : 4
    },
    "eventDetails" : {
      "eventID" : "c5f615f1",
      "customerId" : "123456789",
      "Name" : "NEW"
    }
  }
}

解决方案

方案1:使用JoltTransformJSON处理器

NiFi的JoltTransformJSON处理器适合JSON结构转换,通过Shift操作可直接引用NiFi系统属性:

  1. 把JoltTransformJSON处理器加入流程。
  2. 配置Jolt Specification为Inline Jolt Specification,粘贴以下规范:
[
  {
    "operation": "shift",
    "spec": {
      "body": {
        "metadata": {
          "*": "body.metadata.&",
          "@(kafka.offset)": "body.metadata.kafkaOffset",
          "@(kafka.partition)": "body.metadata.kafkaPartition"
        },
        "eventDetails": "body.eventDetails"
      }
    }
  }
]
  1. 若需要将offset和partition转为整数类型,可追加ModifyOverwriteBeta操作,完整规范如下:
[
  {
    "operation": "shift",
    "spec": {
      "body": {
        "metadata": {
          "*": "body.metadata.&",
          "@(kafka.offset)": "body.metadata.kafkaOffset",
          "@(kafka.partition)": "body.metadata.kafkaPartition"
        },
        "eventDetails": "body.eventDetails"
      }
    }
  },
  {
    "operation": "modify-overwrite-beta",
    "spec": {
      "body": {
        "metadata": {
          "kafkaOffset": "=toInteger",
          "kafkaPartition": "=toInteger"
        }
      }
    }
  }
]

方案2:使用UpdateRecord处理器

如果不熟悉Jolt语法,UpdateRecord处理器操作更直观:

  1. 添加UpdateRecord处理器,配置Record Reader为JsonTreeReader,Record Writer为JsonRecordSetWriter。
  2. 在处理器属性中添加两个更新规则:
    • 键:/body/metadata/kafkaOffset,值:${kafka.offset}
    • 键:/body/metadata/kafkaPartition,值:${kafka.partition}
  3. 运行流程后,消息会自动将NiFi中的Kafka属性写入指定字段。

注意事项

  • NiFi从Kafka消费消息时,会自动携带kafka.offset和kafka.partition系统属性,无需额外配置即可直接引用。
  • 使用UpdateRecord时,需确保Record Reader/Writer的配置与消息JSON格式匹配。

内容的提问来源于stack exchange,提问作者Sam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 18:40:15