如何基于SFTP文件触发/跳过Airflow DAG执行并解决传感器持续连接问题
问题解决:Airflow SFTPSensor文件不存在时持续连接SFTP的问题
需求场景
每5至10分钟调度Airflow DAG检查SFTP服务器指定路径:
- 若存在目标文件(
*.txt),则执行下载、加载等下游任务 - 若文件不存在,直接跳过下游任务
已尝试方案
使用SFTPSensor + BashOperator编写DAG,但遇到问题:文件不存在时,传感器任务会持续连接SFTP服务器,不会终止并跳过下游。
原始代码:
from airflow import DAG from airflow.operators.bash import BashOperator from airflow.providers.sftp.sensors.sftp import SFTPSensor from datetime import date import pandas as pd from future.core.common_tasks import dag_start, dag_end import pendulum import os from airflow.models import Variable dag_id = os.path.basename(__file__).replace(".pyc", "").replace(".py", "") with DAG( dag_id=dag_id, schedule = "*/5 * * * *", is_paused_upon_creation=True, start_date=pendulum.datetime(2024, 1, 1, tz="UTC"), max_active_runs= 1, catchup=False ): wait_for_new_file= SFTPSensor( task_id="wait_for_new_file", path=fullpath, file_pattern='*.txt', sftp_conn_id="sftp_conn", #deferrable=True, newer_than=None ) get_file = BashOperator( task_id='get_sql_tracedata', # 省略下载逻辑 ) load_file = BashOperator( task_id='loadfile', # 省略加载逻辑 ) wait_for_new_file >> get_file >> load_file
问题原因
默认的SFTPSensor采用轮询(poke)模式:它会按照配置的poke_interval(默认30秒)重复连接SFTP服务器检查文件,直到文件出现或达到timeout(默认7天)才会停止。这与你“单次检查、不存在则跳过下游”的需求完全不符。
解决方案
推荐使用ShortCircuitOperator替代SFTPSensor,实现单次检查逻辑,灵活控制下游任务是否执行。
步骤1:导入依赖
添加ShortCircuitOperator和SFTPHook的导入:
from airflow.operators.python import ShortCircuitOperator from airflow.providers.sftp.hooks.sftp import SFTPHook
步骤2:编写SFTP文件检查函数
实现一个Python函数,检查指定SFTP路径下是否存在目标文件:
def check_sftp_file_exists(sftp_conn_id, path, file_pattern): # 初始化SFTP Hook hook = SFTPHook(sftp_conn_id=sftp_conn_id) try: # 获取匹配指定模式的文件列表 matching_files = hook.get_files(path=path, pattern=file_pattern) # 存在匹配文件则返回True,否则返回False return len(matching_files) > 0 finally: # 确保关闭SFTP连接 hook.close()
步骤3:替换SFTPSensor为ShortCircuitOperator
在DAG中替换原有传感器任务:
check_file_exists = ShortCircuitOperator( task_id='check_file_exists', python_callable=check_sftp_file_exists, op_kwargs={ 'sftp_conn_id': 'sftp_conn', 'path': fullpath, 'file_pattern': '*.txt' } )
步骤4:调整任务依赖
保持原有下游任务的依赖关系:
check_file_exists >> get_file >> load_file
效果说明
- 每次DAG调度时,
check_file_exists任务仅连接SFTP一次,检查目标文件是否存在 - 若存在文件,返回
True,下游的下载、加载任务正常执行 - 若不存在文件,返回
False,ShortCircuitOperator会直接跳过所有下游任务,不会持续连接SFTP
内容的提问来源于stack exchange,提问作者Bala Murali
相关产品推荐
相关产品推荐

