如何使用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系统属性:
- 把JoltTransformJSON处理器加入流程。
- 配置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" } } } ]
- 若需要将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处理器操作更直观:
- 添加UpdateRecord处理器,配置Record Reader为
JsonTreeReader,Record Writer为JsonRecordSetWriter。 - 在处理器属性中添加两个更新规则:
- 键:
/body/metadata/kafkaOffset,值:${kafka.offset} - 键:
/body/metadata/kafkaPartition,值:${kafka.partition}
- 键:
- 运行流程后,消息会自动将NiFi中的Kafka属性写入指定字段。
注意事项
- NiFi从Kafka消费消息时,会自动携带
kafka.offset和kafka.partition系统属性,无需额外配置即可直接引用。 - 使用UpdateRecord时,需确保Record Reader/Writer的配置与消息JSON格式匹配。
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

