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

Airflow轮询实现:Sensor的poke_interval与DAG调度配置咨询

如何配置Airflow Sensor实现定期轮询外部系统并避免失败信息

嘿,作为Airflow新手,搞懂Sensor和DAG调度的联动确实容易懵,我来给你拆解清楚,帮你搞定这个场景!

一、先理清核心概念:DAG调度 vs Sensor的poke_interval

很多新手都会混淆这两个参数,先把它们的边界搞明白:

  • DAG的schedule_interval:决定Airflow什么时候启动一个新的DAG Run。比如你设置成@hourly,Airflow就会每小时创建一个DAG实例,尝试执行里面的任务。这是“启动轮询周期”的开关。
  • Sensor的poke_interval:是Sensor本身每隔多久去检查一次外部系统(比如FTP)的条件是否满足。比如设为300秒,就是每5分钟查一次文件是否存在/符合条件。这是“单次轮询内的检查频率”。

简单说:DAG调度是“多久发起一轮检查”,poke_interval是“每轮检查里多久查一次”。

二、针对你的场景的参数配置方案

你的核心需求是:仅在文件可用时触发后续任务,无文件时不产生失败信息。按照这个目标,给你推荐这些配置:

1. DAG层面配置

设置schedule_interval为你想要的轮询周期,比如每小时一次:

schedule_interval="@hourly"  # 或者用 cron 表达式 "0 * * * *"

这个周期决定了Airflow多久会启动一次新的轮询流程。同时记得开启catchup=False,避免Airflow补跑历史DAG Run导致重复检查。

2. Sensor层面关键配置

假设你已经写了自定义Sensor(比如FTPConditionSensor),重点配置这几个参数:

  • mode="reschedule":这是优化资源占用的关键!默认的poke模式会让Sensor一直占用一个Worker槽位,直到条件满足或超时;而reschedule模式会在每次检查后释放槽位,到下一个poke_interval时间再重新申请槽位检查,适合长时间轮询的场景。
  • soft_fail=True:如果Sensor在timeout时间内没找到符合条件的文件,会把这个Sensor任务标记为跳过(Skipped),而不是失败(Failed),这样你的仪表盘就不会被失败信息堆满了。
  • timeout:设置Sensor最长等待时间,建议和你的DAG调度周期匹配。比如DAG每小时调度一次,就设timeout=3600(3600秒=1小时),这样到时间还没找到文件,就直接跳过,等下一轮DAG Run再检查。
  • poke_interval:根据你的需求设置检查间隔,比如5分钟(300秒):poke_interval=300。不要设得太频繁(比如10秒),不然会频繁访问FTP服务器造成压力;也不要太长,避免错过文件。

举个自定义Sensor的配置示例:

from airflow.sensors.base import BaseSensorOperator
from airflow.hooks.ftp_hook import FTPHook
from airflow.utils.decorators import apply_defaults
from datetime import datetime

class FTPConditionSensor(BaseSensorOperator):
    @apply_defaults
    def __init__(self, ftp_conn_id, file_pattern, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.ftp_conn_id = ftp_conn_id
        self.file_pattern = file_pattern

    def poke(self, context):
        # 这里写你的检查逻辑:连接FTP,查找符合file_pattern的新文件
        # 可以额外判断文件修改时间,避免重复处理旧文件
        ftp_hook = FTPHook(ftp_conn_id=self.ftp_conn_id)
        files = ftp_hook.list_directory("/path/to/ftp/target/folder")
        # 示例:匹配文件名包含指定前缀,且是最近1小时内的文件
        matching_files = [
            f for f in files 
            if self.file_pattern in f 
            and ftp_hook.get_mod_time(f) >= datetime.now() - timedelta(hours=1)
        ]
        # 返回True表示找到符合条件的文件,Sensor会触发后续任务;返回False则继续等待
        return len(matching_files) > 0

# 在DAG中使用这个Sensor
with DAG(
    dag_id="ftp_file_processing_workflow",
    schedule_interval="@hourly",
    start_date=datetime(2024, 1, 1),
    catchup=False,
) as dag:
    check_ftp_files = FTPConditionSensor(
        task_id="check_ftp_for_new_files",
        ftp_conn_id="my_ftp_connection",  # 提前在Airflow UI配置好FTP连接
        file_pattern="daily_transaction_",
        poke_interval=300,  # 5分钟查一次
        mode="reschedule",
        soft_fail=True,
        timeout=3600,  # 最长等1小时,和DAG调度周期匹配
    )

    # 后续任务:比如下载文件、解析数据、入库等
    download_ftp_file = BashOperator(
        task_id="download_file_from_ftp",
        bash_command="wget ftp://{{ conn.my_ftp_connection.host }}/path/to/ftp/target/folder/{{ ti.xcom_pull(task_ids='check_ftp_for_new_files') }} -P /local/storage/"
    )

    process_file = PythonOperator(
        task_id="process_downloaded_file",
        python_callable=your_file_processing_function,
        op_kwargs={"file_path": "/local/storage/{{ ti.xcom_pull(task_ids='check_ftp_for_new_files') }}"}
    )

    check_ftp_files >> download_ftp_file >> process_file

三、其他可行方案

如果Sensor的方式不是最贴合你的需求,还有这两个备选:

1. ShortCircuitOperator 替代Sensor

先写一个任务检查FTP是否有符合条件的文件,然后用ShortCircuitOperator根据检查结果决定是否执行后续任务:

def check_ftp_files_func(**context):
    ftp_hook = FTPHook(ftp_conn_id="my_ftp_connection")
    files = ftp_hook.list_directory("/path/to/ftp/folder")
    matching_files = [f for f in files if "daily_transaction_" in f]
    # 把找到的文件名推送到XCom,供后续任务使用
    if matching_files:
        context["ti"].xcom_push(key="target_file", value=matching_files[0])
    return len(matching_files) > 0

with DAG(...) as dag:
    check_files = PythonOperator(
        task_id="check_ftp_files",
        python_callable=check_ftp_files_func,
        provide_context=True
    )

    short_circuit = ShortCircuitOperator(
        task_id="short_circuit_workflow",
        python_callable=lambda context: context["ti"].xcom_pull(task_ids="check_ftp_files")
    )

    download_file = BashOperator(...)

    check_files >> short_circuit >> download_file

这种方式的好处是,没有文件时整个后续流程都会被跳过,不会占用资源;缺点是没有Sensor的“持续检查”能力,只能在DAG启动时检查一次,如果文件是在DAG启动后才上传的,就会错过。

2. 外部触发(非轮询)

如果你的外部系统(比如FTP服务器)支持主动通知,那可以不用轮询,改成事件触发:

  • 当FTP服务器有新文件符合条件时,调用Airflow的REST API触发DAG运行;
  • 或者用消息队列(比如Kafka、RabbitMQ),FTP服务器有新文件时发送消息,Airflow用TriggerDagRunOperator监听消息触发DAG。
    这种方式更高效,完全避免了轮询的资源消耗,但需要外部系统配合实现触发逻辑。

四、避坑提醒

  • 不要把schedule_interval设得太频繁,同时poke_interval又太短,不然会给Airflow和外部系统造成不必要的压力;
  • 处理完文件后,记得添加一个任务把文件移动到归档文件夹,或者标记为已处理,避免后续DAG Run重复处理;
  • 如果你的FTP服务器有访问频率限制,要根据限制调整poke_interval,避免被封禁。

内容的提问来源于stack exchange,提问作者Anton Kupias

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:35:08