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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 20:52:06