Airflow升级后GCSToSFTPOperator通配符功能失效求助
问题分析与解决方案
错误原因
升级后的GCSToSFTPOperator在处理多文件(通配符匹配)且keep_directory_structure=False时,逻辑发生了变化:它不再自动将源文件的文件名追加到目标目录路径后,而是直接尝试将临时文件写入指定的目标目录,导致SFTP客户端抛出OSError: Specified file is a directory错误。
解决方案
方案1:调整参数尝试修复
先尝试修改destination_path的格式,去掉结尾的斜杠,部分新版本的Operator会自动识别无斜杠的路径为目录,并自动追加文件名:
from airflow.providers.google.cloud.transfers.gcs_to_sftp import GCSToSFTPOperator GCSToSFTPOperator( sftp_conn_id=DEFAULT_SFTP_CONN, task_id="export_output_files_to_sftp", destination_path="/FTP_dir/Output", # 移除结尾斜杠 source_bucket="bucket_name", source_object="folder_name/*.ndjson", keep_directory_structure=False, retries=2, move_object=move_object)
方案2:列表+批量复制(可靠替代)
通过GCSListOperator先获取所有匹配的文件,再循环创建GCSToSFTPOperator逐个复制,这种方式逻辑清晰且兼容性强:
from airflow.providers.google.cloud.operators.gcs import GCSListOperator from airflow.providers.google.cloud.transfers.gcs_to_sftp import GCSToSFTPOperator from airflow.operators.python import PythonOperator from airflow.utils.task_group import TaskGroup # 1. 列出GCS中所有.ndjson文件 list_gcs_files = GCSListOperator( task_id="list_gcs_files", bucket="bucket_name", prefix="folder_name/", ) # 2. 循环复制每个文件到SFTP def batch_copy_files(**context): file_list = context['ti'].xcom_pull(task_ids='list_gcs_files') ndjson_files = [f for f in file_list if f.endswith('.ndjson')] with TaskGroup(group_id="copy_files_group") as copy_group: for file_path in ndjson_files: file_name = file_path.split('/')[-1] GCSToSFTPOperator( task_id=f"copy_{file_name}", sftp_conn_id=DEFAULT_SFTP_CONN, source_bucket="bucket_name", source_object=file_path, destination_path=f"/FTP_dir/Output/{file_name}", keep_directory_structure=False, retries=2, move_object=move_object, task_group=copy_group ) return copy_group copy_files_task = PythonOperator( task_id="batch_copy_to_sftp", python_callable=batch_copy_files, provide_context=True, ) list_gcs_files >> copy_files_task
方案3:直接使用Hook手动复制(灵活可控)
如果需要更精细的控制逻辑,可以直接调用GCSHook和SFTPHook完成文件传输:
from airflow.providers.google.cloud.hooks.gcs import GCSHook from airflow.providers.sftp.hooks.sftp import SFTPHook from airflow.operators.python import PythonOperator import tempfile def transfer_files_manual(**context): gcs_hook = GCSHook() sftp_hook = SFTPHook(ftp_conn_id=DEFAULT_SFTP_CONN) # 获取目标文件列表 all_files = gcs_hook.list(bucket_name="bucket_name", prefix="folder_name/") ndjson_files = [f for f in all_files if f.endswith('.ndjson')] for file_path in ndjson_files: file_name = file_path.split('/')[-1] # 下载到本地临时文件 with tempfile.NamedTemporaryFile(mode='wb', delete=False) as tmp_file: gcs_hook.download(bucket_name="bucket_name", object_name=file_path, filename=tmp_file.name) # 上传到SFTP服务器 sftp_dest = f"/FTP_dir/Output/{file_name}" sftp_hook.store_file(sftp_dest, tmp_file.name) # 若需要移动文件(删除GCS源文件) if move_object: gcs_hook.delete(bucket_name="bucket_name", object_name=file_path) manual_transfer_task = PythonOperator( task_id="manual_transfer_to_sftp", python_callable=transfer_files_manual, provide_context=True, )
内容的提问来源于stack exchange,提问作者Martijncl
相关产品推荐
相关产品推荐

