【问题标题】:Read flow file attribute/content to processor property将流文件属性/内容读取到处理器属性
【发布时间】:2021-07-09 05:01:14
【问题描述】:

我想根据通过的最后一个流文件的内容设置处理器的属性。

示例:我使用处理器 GenerateFlowFile 和自定义文本 ${now()} 实例化流文件作为创建流文件期间的当前时间戳。

我想要一个处理器(与我无关)将流文件的内容(时间戳)读取到处理器的自定义属性property_name。之后我希望能够通过 REST-API 查询处理器并从处理器读取该属性。

最初我认为我可以使用ExtractText 处理器来做到这一点,但它会根据正则表达式提取文本并将其写回流文件,而我想将该信息保存在处理器中直到下一个流文件到达。

【问题讨论】:

    标签: apache-nifi flowfile


    【解决方案1】:

    你不能通过 NiFi 做到这一点。当处理器运行时,您无法更新其配置。

    也许您可以在 UpdateAttribute 上使用状态变量?

    有状态使用

    通过为“存储状态”选择“本地存储状态”选项 property UpdateAttribute 不仅会存储评估的属性 作为 FlowFile 的属性,但也作为有状态变量 以递归方式引用。这使处理器能够 计算诸如传入流文件的总和或计数之类的东西。一种 动态属性可以像这样被引用为有状态变量:

    动态属性键:theCount 值: ${getStateValue("theCount"):plus(1)} 这个例子将保持计数 通过处理器的流文件总数。 要在 State 之上使用逻辑,只需使用 更新属性。所有动作都将作为有状态属性存储为 以及被添加到FlowFiles。使用“高级用法”它是 可以跟踪诸如流量最大值之类的事情 远的。这将通过具有以下条件来完成 "${getStateValue("maxValue"):lt(${value})}" 和一个动作 属性:“maxValue”,值:“${value}”。 “状态变量 Initial Value”属性用于初始化有状态变量 如果有状态运行,则需要设置。一些逻辑规则将 需要非常高的初始值,例如使用高级规则 确定最小值。如果有状态属性引用其他 有状态属性,然后是其他有状态属性的值 后面会有一个迭代。例如,试图计算 传入流的平均值需要总和和计数。我摔倒 然后在同一个 UpdateAttribute 中设置三个属性(如下所示) 平均值将始终不包括最新的计数值 和总和:

    计数键:计数值:${getStateValue("theCount"):plus(1)} Sum> key : theSum value : ${getStateValue("theSum"):plus(${flowfileValue})} 平均键:平均值: ${getStateValue("theSum"):divide(getStateValue("theCount"))} 相反, 因为 average 只依赖于 theCount 和 theSum 属性(它们是 也添加到流文件中)应该有以下无状态 UpdateAttribute 正确计算平均值。在活动中 处理器无法在开始时获得状态 onTrigger,FlowFile 将被推回原点 关系和处理器将屈服。如果处理器能够 在 onTrigger 开始时获取状态但无法设置 将属性添加到流文件后的状态,流文件将是 转移到“设置状态失败”。这通常是由于状态不 是最新版本(另一个线程已替换 状态与另一个版本)。在大多数用例中,这种关系 应该循环回处理器,因为唯一受影响的属性 将被覆盖。注意:目前唯一的“有状态”选项是 在本地存储状态。这样做是因为当前的实现 集群状态依赖于 Zookeeper 并且 Zookeeper 不是设计的 对于具有状态的负载/吞吐量 UpdateAttribute 类型 要求。将来,如果/当多个不同的集群状态 添加选项,UpdateAttribute 将被更新。

    【讨论】:

    • 这太棒了!花了我一段时间,但我能够利用theCount 示例为通过处理器的流文件的数量创建一个计数器 - 不完全是时间戳,但我可以跟踪的下一个最好的事情是自从我上次检查以来,数字上升了。最大的好处:我可以通过 REST-API 获得处理器的状态,包括 theCount 的当前值:/processors/{id}/state
    【解决方案2】:

    感谢@Ivan,我能够创建一个完整的工作解决方案 - 以供将来参考:

    1. 使用例如实例化流文件GenerateFlowFile 处理器并添加自定义属性“myproperty”和值 ${now()}(注意:您可以将此属性添加到任何处理器中的流文件中,不必是 GenerateFlowFile 处理器)

    2. 有一个UpdateAttribute 处理器,其选项(在处理器属性下)Store State 设置为Store state locally。

    3. 在UpdateAttribute 处理器中添加一个名为readable_property 的自定义属性,并将其设置为值${'myproperty'}。

    处理器的状态现在包含最后一个流文件的值(例如,带有属性添加到流文件时的时间戳)。

    额外奖励:

    1. 通过 REST-API 和 URI /nifi-api/processors/{id}/state 上的 GET 获取有状态处理器的值(以及通过 (!) 的最后一个流文件的值)

    返回的 JSON 包含以下几行:

    {
    "key":"readable_property"
    ,"value":"Wed Apr 14 11:13:40 CEST 2021"
    ,"clusterNodeId":"some-id-0d8eb6052"
    ,"clusterNodeAddress":"some-host:port-number"
    }
    

    然后你只需要解析 JSON 的值。

    【讨论】:

      【解决方案3】:

      您应该使用 UpdateAttribute 处理器。 您可以阅读几种方法 - f.e. Update attributes based on content in NiFi

      【讨论】:

      • 该链接是一个完全不同的主题 - 我不想在流文件中添加属性...我想更新 处理器的属性 不是流文件。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-02-03
      • 1970-01-01
      • 2019-06-13
      • 1970-01-01
      • 1970-01-01
      • 2019-12-09
      相关资源
      最近更新 更多