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

如何在Apache NiFi中用ExecuteScript的Python脚本获取FTP文件列表

解决Apache NiFi ExecuteScript Python脚本获取FTP文件列表的错误

你的错误核心有两个:

  • 原始流入的flowFile未被指定转移关系或移除,NiFi要求每个流入的FlowFile必须被明确处理(转移或移除);
  • 循环中重复修改同一个newFlowFile,最终仅最后一个文件名对应的FlowFile被转移,且未处理空文件列表的边界情况。

修正后的Python脚本

import ftplib

flowFile = session.get()
if flowFile is not None:
    try:
        # 从FlowFile属性获取FTP服务器配置
        ftp_host = flowFile.getAttribute('HostName')
        ftp_port = int(flowFile.getAttribute('PortNumber'))
        ftp_user = flowFile.getAttribute('Username')
        ftp_password = flowFile.getAttribute('Password')
        ftp_directory = flowFile.getAttribute('RemotePath')

        # 建立FTP连接并登录
        ftp = ftplib.FTP()
        ftp.connect(ftp_host, ftp_port)
        ftp.login(ftp_user, ftp_password)
        ftp.cwd(ftp_directory)

        # 获取目标目录下的文件列表
        file_list = ftp.nlst()

        # 为每个文件名创建独立的FlowFile
        if file_list:
            for filename in file_list:
                # 基于原FlowFile创建新实例(可继承原属性)
                new_flow_file = session.create(flowFile)
                # 设置filename属性
                new_flow_file = session.putAttribute(new_flow_file, 'filename', filename)
                # 将新FlowFile转移到REL_SUCCESS
                session.transfer(new_flow_file, REL_SUCCESS)
        else:
            # 目录为空时,将原FlowFile转移到REL_FAILURE(可根据需求调整逻辑)
            session.transfer(flowFile, REL_FAILURE)

        # 关闭FTP连接
        ftp.quit()
    except Exception as e:
        # 捕获异常,将原FlowFile转移到失败分支并记录日志
        session.transfer(flowFile, REL_FAILURE)
        log.error(f"FTP操作失败: {str(e)}")
    finally:
        # 确保原FlowFile被处理,避免残留
        if flowFile is not None:
            session.remove(flowFile)
else:
    # 无FlowFile流入时直接返回
    pass

关键修正说明

  1. 原FlowFile处理:通过finally块确保原FlowFile被移除,或在分支中明确转移,彻底解决“转移关系未指定”的错误。
  2. 独立FlowFile创建:每次循环创建新的new_flow_file并立即转移,保证每个文件名对应一个独立的FlowFile,不会覆盖之前的实例。
  3. 资源与异常处理:用try-finally保证FTP连接关闭,捕获异常并记录日志,便于问题排查。
  4. 边界情况处理:针对空文件列表场景添加明确逻辑,避免无FlowFile可转移的问题。

可选优化

如果不需要继承原FlowFile的属性,可直接创建全新的FlowFile:

new_flow_file = session.create()

这种方式生成的FlowFile仅包含你设置的filename属性,适合只传递文件名的场景。

内容的提问来源于stack exchange,提问作者Athenkosi Lengs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:53:18