Airflow如何实现任务运行15分钟后再执行下游任务?
原有方案失效原因
你写的time.sleep(900)逻辑从根上就不符合Airflow的调度规则:你配置的DAG调度间隔是1分钟,意味着每分钟都会生成一个独立的DAG Run实例,每个实例里的休眠任务都会单独卡15分钟,既做不到累计15次写入再触发下游,还会长时间占用Worker进程槽位,大概率触发任务运行超时、Worker资源占满的问题,完全达不到预期效果。
可行实现方案
不要在任务代码里写硬休眠,用Airflow原生能力实现,两种方案按需选:
方案1:双DAG+数据集调度(Airflow 2.4及以上版本首选,生产环境推荐)
把写数据和分析数据拆成两个独立DAG,用Airflow的数据集能力做触发联动:
- 第一个DAG保持1分钟调度,只负责每分钟执行写文件逻辑,声明产出对应数据文件的数据集
- 第二个DAG绑定该数据集作为触发条件,配置15分钟的累计窗口,等数据集攒够15分钟的写入更新后,自动触发后续分析任务
参考代码:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.datasets import Dataset from datetime import datetime, timedelta # 定义采集数据文件对应的数据集标识 COLLECTED_DATA = Dataset("file:///path/to/your/collected_data.log") # DAG1:每分钟执行一次数据写入 with DAG( dag_id="collect_data_per_min", schedule_interval="* * * * *", start_date=datetime(2024, 1, 1), catchup=False ) as collect_dag: write_task = PythonOperator( task_id="write_data_to_file", python_callable=your_existing_write_function, # 替换成你实际写文件的业务函数 outlets=[COLLECTED_DATA] ) # DAG2:累计15分钟写入后自动触发分析 with DAG( dag_id="analyze_data_per_15min", schedule=COLLECTED_DATA, start_date=datetime(2024, 1, 1), catchup=False, dataset_condition_timedelta=timedelta(minutes=15) ) as analyze_dag: downstream_task = PythonOperator( task_id="run_analysis", python_callable=your_existing_downstream_function # 替换成你的后续业务任务 )
方案2:单DAG短路判断(适配2.4以下低版本Airflow)
如果版本不支持数据集特性,保持单DAG1分钟调度,在下游任务前加短路判断逻辑:每次DAG运行时先检查是否已经累计采集满15分钟,不满就直接跳过下游任务,满了才放行执行分析,同时重置计数进入下一轮采集周期。
参考代码:
import os import time from airflow import DAG from airflow.operators.python import PythonOperator, ShortCircuitOperator from datetime import datetime TIMESTAMP_MARK_PATH = "/tmp/collect_start_ts.mark" def write_data_with_mark(): # 第一次写入时记录采集起始时间 if not os.path.exists(TIMESTAMP_MARK_PATH): with open(TIMESTAMP_MARK_PATH, "w") as f: f.write(str(time.time())) # 执行原有的每分钟写文件逻辑 your_existing_write_function() def check_collect_duration(): if not os.path.exists(TIMESTAMP_MARK_PATH): return False with open(TIMESTAMP_MARK_PATH, "r") as f: start_ts = float(f.read()) # 累计满15分钟就放行,同时删除标记重置下一轮 if time.time() - start_ts >= 900: os.remove(TIMESTAMP_MARK_PATH) return True return False with DAG( dag_id="collect_and_analyze", schedule_interval="* * * * *", start_date=datetime(2024, 1, 1), catchup=False ) as my_dag: task1 = PythonOperator( task_id="write_data_per_min", python_callable=write_data_with_mark ) check_task = ShortCircuitOperator( task_id="check_if_15min_collected", python_callable=check_collect_duration ) task2 = PythonOperator( task_id="run_downstream_tasks", python_callable=your_existing_downstream_function ) task1 >> check_task >> task2
避坑提醒
不要在Airflow任务里写长耗时sleep逻辑:Airflow Worker的并发槽位是有限的,长睡眠任务会一直占用槽位不释放,很容易把Worker堵死,导致其他任务无法调度。如果业务逻辑允许,也可以直接把DAG调度间隔改成15分钟,在写数据任务内循环15次每分钟完成一次写入再跑下游,但这种方案容错性差,写入中途出错会丢失整轮采集数据,不建议生产环境使用。
内容的提问来源于stack exchange,提问作者oussama benbelhassen
相关产品推荐
相关产品推荐

