如何在Airflow中配置SFTPSensor检测SFTP服务器新增文件
问题分析与解决方案
1. 报错原因
你用的SFTPSensor根本没有file_pattern这个参数,这就是报错的直接原因。这个传感器默认是用来检查单个指定文件是否存在的,本身不支持通过通配符匹配目录下的文件。
2. 解决方案
针对你“监控目录,只要有任意文件就触发后续任务”的需求,给你两种可行方案:
方案一:用支持通配符的SFTPSensor(Airflow 2.3+ 且 apache-airflow-providers-sftp >=3.0.0)
如果你的Airflow环境满足版本要求,可以直接在path里写通配符路径,再开启wildcard=True参数就能实现:
check_SFTP = SFTPSensor( task_id="check_SFTP", sftp_conn_id="My_company", path="/in/*", # 用*匹配目录下所有文件 wildcard=True, # 开启通配符支持 poke_interval=15, timeout=60*5, dag=dag )
方案二:自定义PythonSensor(兼容所有版本)
要是版本不支持方案一,或者需要更灵活的检查逻辑,直接用PythonSensor自己写检查SFTP目录有没有文件的逻辑就行:
def check_sftp_has_files(): ftp_hook = SFTPHook(ftp_conn_id="My_company") files_list = ftp_hook.list_directory("/in/") # 过滤掉.或..这类特殊条目,只留实际文件 real_files = [f for f in files_list if not f.startswith('.')] return len(real_files) > 0 # 替换原来的SFTPSensor check_SFTP = PythonSensor( task_id="check_SFTP", python_callable=check_sftp_has_files, poke_interval=15, timeout=60*5, dag=dag )
3. 修正代码里的另一个大问题
你代码里的files = get_list_of_files()是在DAG解析阶段执行的,不是DAG运行的时候。这会导致每次Airflow刷新DAG定义就去SFTP拉文件列表,而不是等传感器触发后才拿最新的文件。
如果要动态生成后续的处理任务,推荐用Airflow 2.3+支持的动态任务映射,修改后的完整代码示例:
import logging from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator from airflow.providers.sftp.hooks.sftp import SFTPHook from airflow.sensors.python import PythonSensor from airflow.operators.python import PythonOperator from datetime import datetime, timedelta args = { "owner": "My_company", "start_date": datetime(2022,10,17)} def check_sftp_has_files(): ftp_hook = SFTPHook(ftp_conn_id="My_company") files_list = ftp_hook.list_directory("/in/") real_files = [f for f in files_list if not f.startswith('.')] return len(real_files) > 0 def get_sftp_files(**context): ftp_hook = SFTPHook(ftp_conn_id="My_company") files_list = ftp_hook.list_directory("/in/") real_files = [f for f in files_list if not f.startswith('.')] context['ti'].xcom_push(key='sftp_files', value=real_files) return real_files dag = DAG( dag_id = "Checking_SFTP_Server_with_sensor", default_args=args, schedule_interval="0 8 * * *", dagrun_timeout=timedelta(minutes=1), tags=['My_company']) check_SFTP = PythonSensor( task_id="check_SFTP", python_callable=check_sftp_has_files, poke_interval=15, timeout=60*5, dag=dag ) get_files = PythonOperator( task_id="get_sftp_files", python_callable=get_sftp_files, provide_context=True, dag=dag ) start = DummyOperator(task_id='start', dag=dag) # 用动态任务映射生成处理任务 process_order = TriggerDagRunOperator.partial( task_id="process_order", trigger_dag_id="Processing_the_order", dag=dag ).expand( conf=[{"file_name": file} for file in "{{ ti.xcom_pull(key='sftp_files') }}"] ) end = DummyOperator(task_id='end', dag=dag) # 设置任务依赖 check_SFTP >> get_files >> start >> process_order >> end
补充说明
- 动态任务映射通过
partial()和expand()实现,会根据SFTP获取到的文件列表自动生成对应的处理任务。 get_sftp_files任务用XCom传递文件列表,保证是在DAG运行时拿到的最新文件信息,而不是解析阶段的旧数据。
内容的提问来源于stack exchange,提问作者Nadia
相关产品推荐
相关产品推荐

