【问题标题】:what is a fastest way to remove nifi flowfile content?删除 nifi 流文件内容的最快方法是什么?
【发布时间】:2018-11-15 03:37:02
【问题描述】:

我有一个工作流程,我在其中获取 json 文件作为 rest api 的响应。我在一个会话中收到大约 100k 个文件。所有文件的总大小为 15GB。我必须将每个文件保存到文件系统,我正在这样做。在该过程结束时,我必须等待所有文件都存在,然后才能发送成功消息。

一旦我将文件保存在 FS 中,我将调用 notify+wait。但我不再需要流文件中的 15 GB 数据。所以为了释放一些空间,我想到了使用 replaceText 或 ModifyByte 来清除内容。所以通知+等待运行顺利。此过程的总等待时间为 3 小时。

但是在这两种情况下(replaceText 或 ModifyByte)的处理时间都太长了。

您能否建议清除流文件数据的最快方法。我也不需要任何属性。那么这是一种我可以放弃旧流文件并在中途生成kb流文件的方法吗?

我想要的是类似于 generateflowfile 的东西,但是在中间,所以对于我现有的每个流文件,我可以删除旧的流文件,并生成空白流文件以通知并等待。

谢谢

【问题讨论】:

    标签: apache-nifi


    【解决方案1】:

    NiFi 的 Content Repository 和 FlowFile Repository 基于写时复制机制,因此如果您不更改内容或元数据,那么您不一定要在这些处理器上“保留”15GB。

    话虽如此,如果您只需要磁盘上存在此类流文件(而不是内容或元数据),请尝试使用以下 Groovy 脚本执行 ExecuteScript:

    def flowFiles = session.get(1000)
    flowFiles.each {
       session.transfer(session.create(), REL_SUCCESS)
    }
    session.remove(flowFiles)
    

    此脚本一次最多可以抓取 1000 个流文件,并且对于每个流文件,向下游发送一个空流文件。然后它会删除所有原始传入流文件。

    请注意,这(即您的用例)将“破坏”出处/沿袭链,因此如果您的流程出现问题,您将无法分辨哪些流程文件来自哪个父流程文件等. 这个限制是您看不到执行此类功能的完整处理器的原因之一。

    【讨论】:

    • 谢谢,我认为在比较利弊之后,我认为我可以减慢 replaceText 处理器并且不会丢失血统。
    • 你也可以session.create(it) 也不会丢失lineage(而且它还保持与传入流文件相同的属性)
    【解决方案2】:

    如果您需要保留属性、沿袭和元数据,您可以使用以下代码(一次只能抓取 1 个流文件)。唯一改变的是 UUID,但除此之外,所有内容都会保留 - 当然内容除外。

    f = session.get()
    session.transfer(session.create(f), REL_SUCCESS)
    session.remove(f)
    

    【讨论】:

      猜你喜欢
      • 2016-05-14
      • 2014-04-19
      • 2010-09-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-02-13
      相关资源
      最近更新 更多