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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 05:01:26