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

