【问题标题】:Nifi-1.0.0 - synchronization mechanismNifi-1.0.0 - 同步机制
【发布时间】:2016-09-15 13:56:24
【问题描述】:

NiFi 是否有同步机制来知道某事何时完成处理?

我摄取了一些数据,进行了一些处理,在步骤 N-1 我想知道所有数据都已处理,以便继续(最终)步骤 N。

[GetFile / 1000 000 行] ----> [ Proc1 / process step 0 ] -----> [ Proc2 / process step 1 ] .... [ PutSQL / insert into db ] ---> [ Proc 让我知道我已经在表中插入了所有数据] ----> [例如 ProcN / 对数据运行聚合]

【问题讨论】:

    标签: java synchronization data-synchronization apache-nifi


    【解决方案1】:

    NiFi 并没有真正内置在框架中的显式同步功能,但一些处理器具有帮助同步活动的功能。我可以想出几种可能的方法来使您的流程正常工作:

    • 调度 - 您可以在处理器上使用 CRON 调度来调度 GetFile 和稍后的聚合操作,假设操作的持续时间相对可预测。

    • MonitorActivity - MonitorActivity 处理器可以根据队列中的不活动状态触发流文件。您可以使用 PutSQL 的此下游并在插入停止并且应该开始聚合时触发。

    • MergeContent(简单) - MergeContent 处理器可能会将 PutSQL 的结果聚合到触发聚合操作的单个消息中。您必须对 bin 大小和年龄的属性进行试验才能使其正常工作。

    • MergeContent(碎片整理) - MergeContent 有一个碎片整理策略,旨在将较大文件的碎片关联在一起。它需要在流文件上设置特定属性,请参阅文档底部的“读取属性”部分。行为似乎接近您想要的,但设置这些片段属性可能很困难。

    【讨论】:

    • 我实际上扩展了 MergeContent(Defragment) 的行为以实现我想要的,但我很好奇是否有更好的方法来做到这一点。 MonitorActivity 看起来很有趣,我会研究一下。谢谢!
    【解决方案2】:

    我可能会建议您尝试一下。 NiFi 有一个很好的 API,允许您启动和停止处理器。您可以使用 InvokeHTTP 处理器从 NiFi 中调用此 API。这允许您启动 [ ProcN / Run aggregations on data 例如] 并在运行后再次关闭。您必须确保该处理器不会连续运行。所以你的处理器是:

    GetFile / 1000 000 lines] ----> [ Proc1 / process step 0 ] -----> [ Proc2 / process step 1 ] .... [ PutSQL / insert into db ] ---> [ Proc to let me know that I've inserted all the data in the table ] ----> -----> [ InvokeHTTP to start ProcN / Run Aggregates ] --via API call--> [ ProcN / Run aggregates on data for example ] -----> [ InvokeHTTP to stop ProcN / Run Aggregates ]
    

    我们正在研究这种同步请求的方法 - 向远程方回复消息并防止管道中的消息过多。

    【讨论】:

      猜你喜欢
      • 2017-04-24
      • 1970-01-01
      • 2017-03-08
      • 1970-01-01
      • 1970-01-01
      • 2011-08-22
      • 2022-12-03
      • 2017-02-01
      • 1970-01-01
      相关资源
      最近更新 更多