NiFi中如何根据属性/条件将InvokeHttp数据路由到不同处理器
NiFi 动态路由InvokeHttp输出流的配置方法
核心思路是根据判断条件选对应路由处理器,不需要修改InvokeHttp本身的配置,在InvokeHttp下游增加路由逻辑即可,具体配置步骤如下:
前置准备
- 先把InvokeHttp的
response关系(即接口返回的实际数据流)连到后续路由节点,如果你需要拿响应体里的字段做判断,不要直接读内容做路由,先在中间加属性提取处理器:返回格式是JSON就加EvaluateJsonPath,是XML就加EvaluateXPath,是纯文本就加ExtractText,提前把判断要用到的字段值提取成FlowFile属性,后续路由直接读属性判断,性能比解析整个内容高很多。
举个例子:要根据接口返回的biz_type字段路由,就用EvaluateJsonPath新增一个名为biz_type的属性,JSON路径填$.biz_type,运行后每个流文件都会带上这个属性,直接读值判断就行。
路由配置(按判断场景选择)
场景1:基于属性判断(绝大多数场景选这个,最稳定高效)
直接用RouteOnAttribute处理器,配置逻辑非常简单:
- 打开处理器的配置页,在
Properties栏点右上角加号新增路由规则,每个规则对应一个下游分支 - 规则内容直接写NiFi表达式语言即可,给几个常用的判断示例:
- 根据InvokeHttp自带的响应状态码路由:
${http.status.code:equals(200)}对应请求成功分支 - 根据之前提取的业务类型属性路由:
${biz_type:equals('order')}对应订单数据分支 - 模糊匹配属性值路由:
${topic:contains('inventory')}对应库存相关数据分支 - 多条件组合判断:
${http.status.code:equals(200):and(${biz_type:equals('pay')})}对应支付成功数据分支
- 根据InvokeHttp自带的响应状态码路由:
- 所有规则加完后,处理器的
Relationships列表会自动生成每个规则对应的出口,同时自带一个unmatched出口——所有不满足自定义规则的流文件都会走这个出口,一定要把这个出口连到兜底链路(比如打日志、进死信队列、抛异常告警),不然数据会积压在这个处理器里。
场景2:基于流文件内容本身判断
如果判断逻辑没法提前提取成属性,需要直接匹配流内容的特征,就用RouteOnContent处理器:
- 支持正则匹配、字符串包含、字节匹配等模式,比如内容里包含
"retCode":0就走成功分支,包含"retCode":-1就走重试分支 - 如果判断逻辑特别复杂(比如要做数值计算、多字段交叉校验、自定义业务逻辑判断),可以换
ExecuteGroovyScript处理器,写几行简单的Groovy脚本给流文件加个路由标记属性,再连RouteOnAttribute按标记路由即可,灵活度最高。
校验环节
把路由处理器每个规则对应的出口分别连到对应的下游业务处理器,启动流程后可以通过处理器的数据溯源(Data Provenance)功能查看每个流文件的属性和实际路由路径,调整规则直到符合预期即可。
性能提示:高并发场景下尽量用属性判断做路由,不要每次路由都解析整个流文件内容,两者的吞吐量差距能到10倍以上。
内容的提问来源于stack exchange,提问作者Neeraj
相关产品推荐
相关产品推荐

