如何在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
关键修正说明
- 原FlowFile处理:通过
finally块确保原FlowFile被移除,或在分支中明确转移,彻底解决“转移关系未指定”的错误。 - 独立FlowFile创建:每次循环创建新的
new_flow_file并立即转移,保证每个文件名对应一个独立的FlowFile,不会覆盖之前的实例。 - 资源与异常处理:用
try-finally保证FTP连接关闭,捕获异常并记录日志,便于问题排查。 - 边界情况处理:针对空文件列表场景添加明确逻辑,避免无FlowFile可转移的问题。
可选优化
如果不需要继承原FlowFile的属性,可直接创建全新的FlowFile:
new_flow_file = session.create()
这种方式生成的FlowFile仅包含你设置的filename属性,适合只传递文件名的场景。
内容的提问来源于stack exchange,提问作者Athenkosi Lengs
相关产品推荐
相关产品推荐

