You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在NiFi中用ExecuteScript运行Python脚本读取Unpack Content解压数据

问题描述

我尝试用Python脚本读取Unpack Content处理器输出的解压数据,需通过ExecuteScript处理器运行脚本,但存在两个困惑:

  1. 如何用Python代码获取NiFi内解压文件的路径?
  2. 如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 02:32:39