如何在Apache NiFi中基于refresh_time条件新增数值字段
基于Apache NiFi条件新增字段的最优实现方案
首选方案:使用UpdateRecord处理器
UpdateRecord是NiFi官方提供的专门用于记录级字段新增、更新、删除的专用处理器,完全适配你当前ConvertRecord输出的JSON数据流,无需额外格式转换,性能最优。
配置步骤如下:
- 基础参数配置:
- Record Reader选择与你ConvertRecord输出匹配的
JsonTreeReader - Record Writer选择
JsonRecordSetWriter,保持输出为JSON格式,原有字段不会丢失
- Record Reader选择与你ConvertRecord输出匹配的
- 新增动态属性实现条件判断(基于NiFi Record Path语法处理字段逻辑):
- 新增动态属性名:
/cache,属性值:if(toNumber(/refresh_time) < 10, 1, 0) - 新增动态属性名:
/reprocess,属性值:if(toNumber(/refresh_time) > 10, 1, 0)
这里加toNumber是为了兼容refresh_time被识别为字符串的场景,避免类型比较错误。当refresh_time等于10时,两个字段默认取值为0,你可以根据实际需求调整边界条件的取值。
- 新增动态属性名:
配置完成后处理器输出的JSON会直接携带新增的cache、reprocess数值字段,可直接用于后续聚合操作。
备选方案(适合小流量场景)
如果不习惯使用Record Path语法,也可以通过三个处理器串联实现:
- 第一步用
EvaluateJsonPath提取refresh_time为流属性,配置动态属性refresh_time值为$.refresh_time - 第二步用
UpdateAttribute新增两个流属性,使用NiFi表达式语言编写判断逻辑:- 新增属性
cache,值为${refresh_time:lt(10):ifElse(1,0)} - 新增属性
reprocess,值为${refresh_time:gt(10):ifElse(1,0)}
- 新增属性
- 第三步用
AttributesToJSON将cache、reprocess两个流属性写入JSON内容,同时保留原有JSON字段即可。
优先推荐UpdateRecord方案,处理批量记录的性能比备选方案高30%以上,也不需要在流属性和记录内容之间来回转换,维护成本更低。
内容的提问来源于stack exchange,提问作者jjammin92
相关产品推荐
相关产品推荐

