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

Apache Airflow多文件SFTP传输实现方案咨询

解决方案:Airflow批量传输目录下所有文件至SFTP服务器

针对你的需求,提供两种实现方式,分别适配不同场景:

方式一:单任务内批量传输所有文件

适合文件数量不多、无需单独监控每个文件传输状态的场景,通过PythonOperator调用SFTP Hook完成批量操作:

import os
from airflow import DAG
from datetime import datetime
from airflow.operators.python import PythonOperator
from airflow.providers.sftp.hooks.sftp import SFTPHook

# 替换为你的实际配置
SFTP_CONNECTION_ID = "your_sftp_conn_id"
SFTP_SOURCE_TR_CURL_PATH = "/path/to/local/files/"
SFTP_DESTINATION_TR_CURL_PATH = "/path/to/remote/files/"

conf = {
    "start_date": datetime(2022, 6, 1),
    "catchup": False,
    "schedule_interval": "@daily",
    "dag_id": "KPO_batch_sftp_transfer"
}

def batch_transfer_files():
    # 获取本地目录下所有文件(不含子目录)
    local_files = [f for f in os.listdir(SFTP_SOURCE_TR_CURL_PATH) 
                   if os.path.isfile(os.path.join(SFTP_SOURCE_TR_CURL_PATH, f))]
    
    if not local_files:
        print("本地目录无文件可传输")
        return
    
    # 初始化SFTP Hook
    sftp_hook = SFTPHook(ssh_conn_id=SFTP_CONNECTION_ID)
    
    try:
        # 遍历文件逐一传输
        for file_name in local_files:
            local_path = os.path.join(SFTP_SOURCE_TR_CURL_PATH, file_name)
            remote_path = os.path.join(SFTP_DESTINATION_TR_CURL_PATH, file_name)
            
            print(f"开始传输文件: {file_name}")
            sftp_hook.store_file(remote_path, local_path)
            print(f"文件 {file_name} 传输完成")
    finally:
        # 确保关闭SFTP连接
        sftp_hook.close_conn()

with DAG(**conf) as dag:
    task_batch_transfer = PythonOperator(
        task_id="batch_transfer_all_files",
        python_callable=batch_transfer_files
    )

task_batch_transfer

方式二:动态生成单个文件传输任务

适合需要单独监控每个文件传输状态的场景,遍历文件后为每个文件生成独立的SFTPOperator:

import os
from airflow import DAG
from datetime import datetime
from airflow.providers.sftp.operators.sftp import SFTPOperator, SFTPOperation

# 替换为你的实际配置
SFTP_CONNECTION_ID = "your_sftp_conn_id"
SFTP_SOURCE_TR_CURL_PATH = "/path/to/local/files/"
SFTP_DESTINATION_TR_CURL_PATH = "/path/to/remote/files/"

conf = {
    "start_date": datetime(2022, 6, 1),
    "catchup": False,
    "schedule_interval": "@daily",
    "dag_id": "KPO_dynamic_sftp_transfer"
}

with DAG(**conf) as dag:
    # 获取本地目录下所有文件(不含子目录)
    local_files = [f for f in os.listdir(SFTP_SOURCE_TR_CURL_PATH) 
                   if os.path.isfile(os.path.join(SFTP_SOURCE_TR_CURL_PATH, f))]
    
    if not local_files:
        print("本地目录无文件可传输")
    else:
        # 为每个文件生成独立的传输任务
        for file_name in local_files:
            SFTPOperator(
                task_id=f'put_{file_name}',
                ssh_conn_id=SFTP_CONNECTION_ID,
                local_filepath=os.path.join(SFTP_SOURCE_TR_CURL_PATH, file_name),
                remote_filepath=os.path.join(SFTP_DESTINATION_TR_CURL_PATH, file_name),
                operation=SFTPOperation.PUT,
                create_intermediate_dirs=True
            )

注意事项

  • 包含子目录传输:如果需要传输子目录下的文件,可改用os.walk()遍历,示例:
    local_files = []
    for root, dirs, files in os.walk(SFTP_SOURCE_TR_CURL_PATH):
        for file in files:
            relative_path = os.path.relpath(os.path.join(root, file), SFTP_SOURCE_TR_CURL_PATH)
            local_files.append(relative_path)
    
    此时远程路径需对应创建子目录,可通过os.path.join(SFTP_DESTINATION_TR_CURL_PATH, relative_path)实现。
  • 空目录处理:代码中加入空目录判断,避免无文件时任务无意义执行。
  • 异常扩展:可根据需求在批量传输代码中添加异常捕获、重试逻辑,提升任务稳定性。
  • 变量替换:务必将代码中的配置变量替换为你的实际SFTP连接ID、本地/远程路径。

内容的提问来源于stack exchange,提问作者rf guy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 18:07:13