【问题标题】:Flowfiles stuck in queue (Apache NiFi)流文件卡在队列中(Apache NiFi)
【发布时间】:2021-02-25 13:41:16
【问题描述】:

我有以下流程:

ListFTP -> RouteOnAttribute -> FetchFTP -> UnpackContent -> ExecuteScript.

部分文件卡在队列UnpackContent -> ExecuteScript

ExecuteScript 吃了一些流文件,它们就消失了:failuresuccess 关系为空。它只是在Tasks/Time 字段中显示了一些活动。他们都在ExecuteScript 之前排队。我试图清空队列,但并非所有流文件都已从此队列中删除。其中大约 1/3 仍然排在队列中。我试图再次禁用所有处理器并清空队列,但它返回:0 FlowFiles (0 bytes) were removed from the queue.

当我尝试更改连接目标时,它会返回:

Cannot change destination of Connection because FlowFiles from this Connection are currently held by ExecuteScript[id=d33c9b73-0177-1000-5151-83b7b938de39]

ExecuScript 来自 answer(使用 Python)。

所以,我不能清空队列,因为它总是返回没有任何流文件的消息,而且我不能删除连接。这已经持续了几个小时。

连接配置:

调度设置为 0 秒,流文件等没有惩罚。

是脚本问题吗?

更新

将脚本更改为:

flowFile = session.get() 
if (flowFile != None):
    # All processing code starts at this indent
    if errorOccurred:
        session.transfer(flowFile, REL_FAILURE)
    else:
        session.transfer(flowFile, REL_SUCCESS)
# implicit return at the end

同样的结果。

更新 v2

我将并发任务设置为 50,然后再次运行 ExecuteScript 并终止它。我收到了这个错误:

更新 v3

我用相同的脚本创建了额外的 ExecuteScript 处理器,它工作正常。但是在我停止这个新处理器并创建新的流文件之后,这个处理器现在有同样的问题:它只是卡住了。

搞笑。 ExecuteScript 是一次性使用吗?

【问题讨论】:

  • 听起来你的脚本仍在处理文件
  • 我猜出了点问题。因为它已经处理了 8 个小时的文件。
  • 我用相同的脚本创建了额外的ExecuteScript 处理器,它工作正常。我真的不知道其他ExecuteScript 处理器发生了什么。如何删除这些流文件...我什至无法终止它,因为处理器中有 0 个流文件。
  • 你试过重启 nifi 之类的明显的东西吗?
  • 我现在无法重新启动它,只能在晚上重新启动,因为我们有许多正在工作的(现在)其他处理器。

标签: apache-nifi


【解决方案1】:

您需要修改 nifi-1.13.2 中的代码,因为 NIFI-8080 导致了这些错误。或者你只使用 nifi 1.12.1

JythonScriptEngineConfigurator:

@Override
public Object init(ScriptEngine engine, String scriptBody, String[] modulePaths) throws ScriptException {
    // Always compile when first run
    if (engine != null) {
            // Add prefix for import sys and all jython modules
            prefix = "import sys\n"
                    + Arrays.stream(modulePaths).map((modulePath) -> "sys.path.append(" + PyString.encode_UnicodeEscape(modulePath, true) + ")")
                    .collect(Collectors.joining("\n"));
    }
    return null;
}

@Override
public Object eval(ScriptEngine engine, String scriptBody, String[] modulePaths) throws ScriptException {
    Object returnValue = null;
    if (engine != null) {
        returnValue = ((Compilable) engine).compile(prefix + scriptBody).eval();
    }
    return returnValue;
}

【讨论】:

  • 在哪里可以看到更多信息?我无法阅读小图像中的文字。此代码位于何处?是否需要重新编译,或者是否可以停止 Nifi,进行更改,然后重新启动?
猜你喜欢
  • 1970-01-01
  • 2017-10-03
  • 2020-10-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多