Apache NiFi创建循环解决销售订单重复注册问题
解决NiFi订单重复插入问题的方案
先看你提供的订单数据:
| sales_order | item_code | description |
|---|---|---|
| 100 | 1 | Item 1 |
| 100 | 2 | item 2 |
| 100 | 3 | Item 3 |
| 100 | 4 | Item 4 |
你的问题核心是同一订单的多个明细流文件同时处理,导致主表重复插入。不建议用循环(不符合NiFi流处理的设计逻辑,还容易出问题),推荐两种更高效的方案:
方案一:按订单号聚合流文件(最优)
把同一订单的所有明细合并成一个流文件,统一处理主表插入和明细批量写入:
- 步骤1:聚合流文件
使用MergeContent处理器,配置Attribute to Use for Binning为sales_order,将同一订单的所有明细流文件合并成一个。这样每个订单只会生成一个待处理的流文件。 - 步骤2:检查订单存在性
用LookupRecord或ExecuteSQL查询销售订单表,判断当前sales_order是否已存在:- 若存在:通过
RouteOnAttribute路由到Discard处理器直接丢弃。 - 若不存在:进入下一步。
- 若存在:通过
- 步骤3:写入主表和明细表
- 用
PutSQL执行一次主表插入操作(因为聚合后一个订单只有一个流文件,不会重复插入)。 - 把聚合后的流文件拆分为单个明细(用
SplitRecord),再用PutRecord开启批量写入模式,一次性插入所有明细行;或者直接用支持批量插入的JDBC写入器完成明细插入。
- 用
方案二:用分布式缓存做订单锁
如果不想聚合流文件,可通过分布式缓存标记已处理的订单:
- 步骤1:查询缓存
使用DistributedMapCacheLookup查询缓存中是否存在当前sales_order。 - 步骤2:分支处理
- 若缓存中不存在:先执行
PutSQL插入主表,再用DistributedMapCachePut将该sales_order存入缓存(可设置合理过期时间,避免内存浪费),最后插入明细表。 - 若缓存中已存在:直接插入明细表即可。
- 若缓存中不存在:先执行
为什么不推荐循环?
NiFi是基于流的处理框架,循环会增加流程复杂度,还可能引发死循环、资源占用过高的问题。上述两种方案更贴合NiFi的设计模式,能高效解决重复插入问题。
内容的提问来源于stack exchange,提问作者José Roberto Miesbach
相关产品推荐
相关产品推荐

