【问题标题】:How to pass flow files to the Execute Python script and use attributes & Nifi variables to store that file?如何将流文件传递给执行 Python 脚本并使用属性和 Nifi 变量来存储该文件?
【发布时间】:2019-05-10 14:08:36
【问题描述】:

我是 NiFi 和 Python 的新手,我需要您的帮助才能将 Flow File 属性值传递给脚本。该脚本正在将嵌套的 json 转换为 csv。当我在本地运行脚本时,它可以工作。

如何将 FlowFile 名称传递给 src_json 和 tgt_csv?

谢谢,

罗莎

import pandas as pd
import json
from pandas.io.json import json_normalize

src_json = "C:/Users/name/Documents/Filename.json"
tgt_csv = "C:/Users/name/Documents/Filename.csv"

jfile = open(src_json)
jdata = json.load(jfile)

...rest of the code...
```python

【问题讨论】:

  • 您可以使用ExecuteStreamCommand 来运行您的python 脚本。并重新编写脚本以从标准输入读取 json 并将 csv 写入标准输出。
  • 当'ConvertRecord'之类的处理器随时可以完成这项工作时,我可以知道为什么需要自定义python脚本来将json转换为csv吗?只是好奇
  • @Arun211 你建议如何解决多个 json 模式?我尝试了 JOLT,但我无法获得正确的输出。目前,我只有一个模式结构进来,但将来我会有很多。参考:stackoverflow.com/questions/56061222/…

标签: python json csv apache-nifi


【解决方案1】:

您有几个选择来完成这项任务。

  1. 正如Arun211 所指出的,现有的ConvertRecord 处理器在很大程度上完成了这项任务。如果您的嵌套 JSON 存在问题,或者您有其他原因想要在 Python 脚本中执行此操作,请继续下面的操作。
  2. 如果您有执行上述任务的现有 Python 脚本,则需要在向脚本提供数据时从 NiFi 调用它。您可以使用:
    1. ExecuteScript(更适合原型设计)和InvokeScriptedProcessor(更适合生产任务)允许您在 NiFi 实例中运行 Python(实际上是 Jython)脚本。这使您可以直接访问一些便捷的方法和功能。但是,由于 Jython 无法处理本机编译的 Python 库,因此您将无法在此代码中使用 pandasSee here for instructions on configuring this processorhere for why pandas will not work
    2. 如果您需要 pandas 的某些功能,您需要将脚本保存为本地文件系统上的 Python 文件,并使用 ExecuteStreamCommand 将其作为 shell 命令调用(如果您需要为此处理器提供输入)或ExecuteProcess(如果它是您的流程中的第一个处理器)。这些处理器本质上运行一个 shell 命令,如python my_python_script_with_pandas.py -someargExecuteProcess)或python my_python_script_with_pandas.py,流文件内容为STDINExecuteStreamCommand),STDOUT 的输出被捕获为结果流文件内容。

目前,您的脚本正在静态文件位置查找传入的 JSON 文件,并将生成的 CSV 放在另一个静态文件位置。您需要更改脚本以执行以下操作之一:

  1. 从命令行参数读取这些路径,并在您选择的处理器的相关处理器属性中传递这些路径。这些属性可以从 flowfile 属性 中填充,因此您可以执行类似 Command Arguments:-inputfile /path/to/some_existing_file.json -outputfile ${flowfile_attribute_named_output_file} 或其任意组合的操作。然后,您的脚本将读取 -inputfile-outputfile 参数以确定路径。
  2. 直接从STDINexample here读取传入数据。然后处理 JSON 数据,将其转换为 CSV,并通过STDOUT 返回。 NiFi 将使用这些数据,将其作为结果流文件的内容,并将其发送到流中的下一个处理器。
  3. 前两个选项使您的 Python 脚本独立于 NiFi;它不知道任何“流文件”结构。此选项将使其特定于 NiFi,但允许更多功能(请参阅上面的选项 2.1)。要编写直接读写流文件内容的 Python 代码,请参阅example of ExecuteScript processor handling flowfile content in Python

【讨论】:

  • 您建议如何处理多个传入的 JSON 模式?
  • 这样吗? snag.gy/x4f1yv.jpg 当你说前两个让 Nifi 不知道任何“流文件”结构时,我很困惑。那么,我不能将 source_path_filename 属性传递给 ExecuteStreamCommand?
  • 看起来以这种方式传递属性很好,但您还需要设置 Command Path 以调用直接调用 Python 的 shell 脚本,或者在参数中也提供 Python 脚本名称,并调用 python 作为命令路径。
猜你喜欢
  • 2019-12-12
  • 2020-09-12
  • 2014-05-12
  • 2022-12-01
  • 2023-02-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-09-18
相关资源
最近更新 更多