如何在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
相关产品推荐
相关产品推荐

