如何在NiFi中用ExecuteScript运行Python脚本读取Unpack Content解压数据
问题描述
我尝试用Python脚本读取Unpack Content处理器输出的解压数据,需通过ExecuteScript处理器运行脚本,但存在两个困惑:
- 如何用Python代码获取NiFi内解压文件的路径?
- 如何用Python代码设置NiFi内处理后文件的输出路径?
当前NiFi数据流包含Unpack Content处理器,后续连接ExecuteScript处理器执行Python脚本处理。
现有Python代码(注释已翻译)
#!/usr/bin/env python # 编码格式:utf-8 from sysconfig import get_python_version import pandas as pd import os from pathlib import Path get_python_version().run_line_magic('pip', 'install pyarrow') # 指定解压后insights文件的路径 path = Path(r'C:\Users\IT Admin\Desktop\Parquet_to_CSV_V2\insights_2023-06-29') os.chdir(path) cwd = Path.cwd() cwd # 将"complex_relations"文件夹下所有parquet文件的路径添加到列表 target_dir = cwd / "complex_relation" pq_files = [] for file in target_dir.rglob("*.parquet*"): pq_files.append(file) # 循环读取所有parquet文件并合并成一个DataFrame data_frames=[] for parquet in pq_files: df = pd.read_parquet(parquet) data_frames.append(df) concatenated_df = pd.concat(data_frames) # 指定合并后parquet文件的输出路径 output_path = r'C:\Users\IT Admin\Desktop\Parquet_to_CSV_V2\output\complex_relations.parquet' concatenated_df.to_parquet(output_path, engine = 'pyarrow') compliled_pq = r'C:\Users\IT Admin\Desktop\Parquet_to_CSV_V2\output\complex_relations.parquet' pd.read_parquet(compliled_pq, engine = "auto")
解决方案
1. 获取NiFi内解压文件的路径
NiFi中Unpack Content处理器解压文件夹后,会将解压目录路径写入FlowFile的absolute.path属性。在ExecuteScript的Python脚本中,需通过NiFi提供的flowFile对象读取该属性,不能硬编码本地路径。
示例代码片段:
from org.apache.nifi.processor.io import StreamCallback import os class PyStreamCallback(StreamCallback): def process(self, inputStream, outputStream): # 从FlowFile属性中获取解压后的目录路径 unpacked_dir = flowFile.getAttribute('absolute.path') # 定位到目标parquet文件夹 target_dir = os.path.join(unpacked_dir, "complex_relation") # 后续遍历parquet文件的逻辑...
2. 设置NiFi内处理后文件的输出路径
NiFi的设计逻辑是通过FlowFile传递数据,无需在脚本中硬编码输出路径。正确做法是将处理后的内容写入FlowFile的输出流,再通过后续的PutFile处理器配置最终输出目录。如果需要自定义输出文件名,可在脚本中修改FlowFile的filename属性。
完整适配NiFi的示例脚本:
from org.apache.nifi.processor.io import StreamCallback import java.io import pandas as pd import os import io class PyStreamCallback(StreamCallback): def process(self, inputStream, outputStream): # 获取解压目录 unpacked_dir = flowFile.getAttribute('absolute.path') target_dir = os.path.join(unpacked_dir, "complex_relation") # 收集所有parquet文件路径 pq_files = [] for root, _, files in os.walk(target_dir): for file in files: if file.endswith('.parquet'): pq_files.append(os.path.join(root, file)) # 合并parquet文件 data_frames = [pd.read_parquet(file) for file in pq_files] concatenated_df = pd.concat(data_frames, ignore_index=True) # 将合并结果写入FlowFile输出流 buffer = io.BytesIO() concatenated_df.to_parquet(buffer, engine='pyarrow') buffer.seek(0) outputStream.write(buffer.read()) # 执行回调处理并传递FlowFile flowFile = session.write(flowFile, PyStreamCallback()) # 设置输出文件名 flowFile = session.putAttribute(flowFile, 'filename', 'complex_relations.parquet') session.transfer(flowFile, REL_SUCCESS)
注意:脚本中不能用get_python_version().run_line_magic('pip', 'install pyarrow')安装依赖,必须提前在NiFi使用的Python环境中手动安装好pandas和pyarrow。
内容的提问来源于stack exchange,提问作者mar7y_k
相关产品推荐
相关产品推荐

