在NiFi中,我有一个cron驱动的处理器序列,它每天提供一组流文件,其中包含我感兴趣的两个属性:
product_code
和
publication_date
。
我需要的是每个
product\u代码
:最新的
发布\u日期
。
例如:
对于此输入:
flow_1: product_code: A / publication_date : 2018-01-01
flow_2: product_code: B / publication_date : 2018-01-01
flow_3: product_code: C / publication_date : 2018-01-01
flow_4: product_code: A / publication_date : 2018-04-12
flow_5: product_code: A / publication_date : 2000-12-31
flow_6: product_code: B / publication_date : 2018-02-02
flow_7: product_code: B / publication_date : 2018-03-03
预期输出应为:
flow_3: product_code: C / publication_date : 2018-01-01
flow_4: product_code: A / publication_date : 2018-04-12
flow_7: product_code: B / publication_date : 2018-03-03
我测试的算法
-
使用
UpdateAttribute
添加属性的处理器
priority
到每个流文件,基于
发布\u日期
。
-
这些更新的流文件被重定向到
PriorityAttributePrioritizer
队列
-
流文件保留在此队列中,因为只有一个使用cron驱动的处理器。通过这种方式,我确信队列中的流文件是根据
发布\u日期
。
-
然后CRON触发下一个处理器
DetectDuplicate
基于
product\u代码
属性由于流文件是从最新的项目处理到最旧的项目,我确信当
product\u代码
被检测为重复,这是因为
product\u代码
对于最近的
发布\u日期
。
问题
可悲的是,当cron触发
检测到重复
处理器,仅使用一条消息,其他消息留在队列中。
如果我将“调度策略”更改为“计时器驱动”,并且“运行调度”为0,那么我的所有流文件都将被消耗,并且输出符合预期。
有没有办法问我
检测到重复
处理器在队列开始工作时使用队列中的所有消息(而不仅仅是一条消息)?
或者有没有一种方法可以设置一种调度策略,比如“凌晨2:00开始工作,凌晨4:00停止”?
你认为有更好的策略来满足需求吗?
当做
Val。
更新1
(2018-04-13)更多信息,除了布莱恩·本德的评论。
我知道CRON不是最好的解决方案,但我不知道如何改进我的算法来摆脱它。
在我的例子中,排队等待重复数据消除的流文件是通过3个REST调用序列生成的:
-
第一次调用“GetAllCategories”,
-
然后,对于每个类别,调用“GetSubCategories”,
-
对于每个子类别,调用“GetProducts”。
此流文件生成部分通常持续5分钟左右:昨晚,第一个流文件在凌晨2:00:16到达队列,最后一个流文件在凌晨2:04:58到达队列。(这就是我安排
检测到重复
在凌晨3:00运行。)
如果我的
检测到重复
处理器将被“计时器驱动”调度,到达队列的第一个流文件将在所有流文件到达之前由处理器消耗。
这将打破全套流文件的顺序。
我觉得我必须等待所有流文件在
检测到重复
处理器开始工作。
您是否有改进我的算法的潜在建议?