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

如何在Airflow中并行执行SFTP文件下载任务(限10并发)

嗨,刚接触Airflow就能考虑到并行优化这点真的很赞!你现在在单个Operator里串行下载的方式,确实没法利用Airflow的多Worker集群能力,不过Airflow有现成的方案完美解决你的需求,下面给你两种最实用的实现思路:

方法一:Dynamic Task Mapping(Airflow 2.3+ 推荐)

这是Airflow官方主推的动态任务生成方式,代码简洁还能轻松控制并发数,完全适配你的场景:

步骤1:先配置SFTP连接

先去Airflow的「Admin -> Connections」页面,新建一个SFTP类型的连接,填好你的SFTP服务器地址、端口、用户名/密码(或密钥),记好连接ID(比如叫sftp_default)。

步骤2:编写DAG代码

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

# 配置你的路径和连接ID
REMOTE_PATH = "/your/remote/dir/"
LOCAL_PATH = "/your/local/dir/"
SFTP_CONN_ID = "sftp_default"

def get_remote_file_list(**context):
    """获取SFTP服务器上需要下载的文件列表"""
    sftp_hook = SFTPHook(sftp_conn_id=SFTP_CONN_ID)
    # 列出远程目录下的所有文件
    all_files = sftp_hook.list_directory(REMOTE_PATH)
    # 过滤出你需要的test_*.csv文件
    target_files = [f for f in all_files if f.startswith("test_") and f.endswith(".csv")]
    # 把文件列表通过XCom传递给下一个任务
    context["ti"].xcom_push(key="remote_files", value=target_files)

def download_single_file(file_name, **context):
    """下载单个文件的函数"""
    sftp_hook = SFTPHook(sftp_conn_id=SFTP_CONN_ID)
    remote_file = f"{REMOTE_PATH}{file_name}"
    local_file = f"{LOCAL_PATH}{file_name}"
    sftp_hook.get(remote_file, local_file)
    print(f"✅ 完成下载: {file_name}")

with DAG(
    dag_id="sftp_parallel_download_dag",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # 按需触发,也可以设成定时
    catchup=False,
    max_active_tis_per_dagrun=10,  # 全局限制DAG最多同时跑10个任务
) as dag:

    # 第一步:获取需要下载的文件列表
    get_files_task = PythonOperator(
        task_id="fetch_remote_file_list",
        python_callable=get_remote_file_list,
        provide_context=True,
    )

    # 第二步:动态生成下载任务,每个文件对应一个任务实例
    download_tasks = PythonOperator(
        task_id="download_single_file",
        python_callable=download_single_file,
        op_kwargs={"file_name": "{{ ti.xcom_pull(key='remote_files') }}"},
        provide_context=True,
        max_active_tis_per_task=10,  # 单独限制这个任务的并发数
    )

    # 设置任务依赖
    get_files_task >> download_tasks

为什么这个方案好用?

  • 每个文件的下载都是独立的Task实例,Airflow会自动把这些任务分配到不同的Worker执行,真正实现并行
  • 通过max_active_tis_per_task或max_active_tis_per_dagrun可以精准控制最多10个并发下载
  • 代码不需要手动写循环生成任务,Airflow会根据文件列表自动映射生成

方法二:循环生成独立Task(兼容Airflow 2.0+)

如果你的Airflow版本低于2.3,不支持Dynamic Mapping,可以用这种手动循环生成任务的方式:

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

REMOTE_PATH = "/your/remote/dir/"
LOCAL_PATH = "/your/local/dir/"
SFTP_CONN_ID = "sftp_default"

def download_single_file(file_name):
    sftp_hook = SFTPHook(sftp_conn_id=SFTP_CONN_ID)
    remote_file = f"{REMOTE_PATH}{file_name}"
    local_file = f"{LOCAL_PATH}{file_name}"
    sftp_hook.get(remote_file, local_file)

with DAG(
    dag_id="sftp_parallel_download_legacy",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False,
    max_active_tis_per_dagrun=10,  # 限制并发数为10
) as dag:

    # 先手动定义你要下载的文件列表(或者用PythonOperator先获取远程列表)
    target_files = [f"test_{i}.csv" for i in range(1, 11)]

    # 循环生成每个文件的下载任务
    for file in target_files:
        download_task = PythonOperator(
            task_id=f"download_{file}",
            python_callable=download_single_file,
            op_kwargs={"file_name": file},
        )
        # 如果有前置任务(比如获取文件列表),可以在这里设置依赖

注意事项

  • 不管用哪种方法,都要确保你的Worker节点能访问到本地下载路径LOCAL_PATH
  • 如果SFTP服务器有连接限制,并发数不要设得比服务器允许的连接数高,避免被限流

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:09:13