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
相关产品推荐
相关产品推荐

