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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 00:05:04