【问题标题】:how to take flowfile count in nifi queue?如何在 nifi 队列中获取流文件计数?
【发布时间】:2017-10-03 13:40:48
【问题描述】:

我有类似的 nifi 流(独立)

executestreamprocessor(hive script) -> executestreamprocessor(hadoop script).

对于每个传入的流文件,hive 脚本使用命令 INSERT..INTO..SELECT..FROM 运行,hadoop 脚本从存储区域中删除特定文件。

有时,当 hadoop 脚本同时运行命令时,hive 脚本会失败。

我每小时最多可以获得 4 个文件。所以我计划在 hive 和 hadoop 处理器之间使用 controlrate 处理器。我设置了队列数量达到4个流文件时的条件,然后应该执行hadoop脚本。但是,controlrate 具有仅设置最大速率的属性。它没有最低费率。

是否有任何可能的解决方案来实现?或任何其他解决方案?

【问题讨论】:

    标签: apache-nifi


    【解决方案1】:

    您应该可以为此使用ExecuteScript,试试这个 Groovy 脚本:

    def flowFiles = session.get(4)
    if(!flowFiles || flowFiles.size() < 4) {
      session.rollback()
    } else {
      session.transfer(flowFiles, REL_SUCCESS)
    }
    

    如果您只想触发一次下游流,那么您可以从父级创建一个子流文件(并报告一个 JOIN 出处事件):

    def flowFiles = session.get(4)
    if(!flowFiles || flowFiles.size() < 4) {
      session.rollback()
    } else {
      def flowFile = session.create(flowFiles)
      session.provenanceReporter.join(flowFiles, flowFile)
      session.remove(flowFiles)
      session.transfer(flowFile, REL_SUCCESS)
    } 
    

    话虽如此,如果您不关心流文件内容(即,您使用流文件作为触发器),您可以使用 MergeContent,最小和最大条目数 = 4。

    【讨论】:

    • 我将在一个位置获得 4 个流文件。我将把这些文件中的每一个都路由到配置单元脚本(executestreamprocessor)和hadoop脚本(executestreamprocessor)。基本上,我将通过 hive 脚本和 hadoop 脚本处理每个文件以删除所有文件。但我不想单独删除文件。在这里,我认为 executescript(groovy) 和 mergecontent 不适合。请多多指教。
    • 脚本是否已经知道要删除哪些文件(就像文件夹中的所有文件一样)?如果是这样,第二个脚本应该可以工作,把它放在 hive 脚本之后和 Hadoop 执行流命令之前,然后你会得到 4 个 Hive 命令,然后是一个 Hadoop 删除
    • 是否可以检查其他关系的第二个脚本的条件?
    • 我不明白你的意思,你是说脚本中的错误处理吗?
    • 假设,我有 1.executestreamcommand(运行 hive 脚本),然后是 2.executestreamcommand(运行 hive 脚本),然后是 3.executestreamcommand(运行 hive 脚本),然后是 4.executestreamcommand(运行 hadoop 命令) .我希望仅当第一个执行流命令队列为空时才执行 hadoop 命令。在这种情况下,我们需要检查第一个处理器的队列。有可能吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多