如何在Airflow的DAG中实现wait_for_dag2步骤,待DAG2完成后执行clean任务?
实现wait_for_dag2步骤的两种方案
方案一:使用PythonSensor自定义检查逻辑
这种方式可以同时满足「dag2无运行中实例」和「dag2最新实例为成功状态」两个条件,支持周期性检查直到条件满足。
- 导入所需模块
from airflow import DAG from airflow.sensors.python import PythonSensor from airflow.models import DagRun from airflow.utils.state import DagRunState from datetime import datetime
- 编写检查dag2状态的poke函数
该函数每次检查时返回布尔值:True表示条件满足,传感器任务完成;False则继续等待下一次检查。
def check_dag2_valid(**context): # 按执行时间倒序获取dag2的所有运行实例 dag2_runs = DagRun.find(dag_id="dag2", order_by=DagRun.execution_date.desc()) # 检查是否存在运行中的dag2实例 has_running = any(run.state == DagRunState.RUNNING for run in dag2_runs) if has_running: return False # 检查dag2最新实例是否为成功状态 if not dag2_runs or dag2_runs[0].state != DagRunState.SUCCESS: return False return True
- 在dag1中定义wait_for_dag2任务
with DAG( dag_id="dag1", start_date=datetime(2024, 1, 1), schedule_interval="@daily" ) as dag: # 原有的start、clean、end任务定义 start = ... clean = ... end = ... # 定义wait_for_dag2传感器任务 wait_for_dag2 = PythonSensor( task_id="wait_for_dag2", python_callable=check_dag2_valid, poke_interval=60, # 每隔60秒检查一次,可按需调整 timeout=3600, # 超时时间(秒),超过则任务失败 mode="reschedule" # 采用reschedule模式节省资源,适合长时间等待场景 ) # 组装任务流 start >> wait_for_dag2 >> clean >> end
方案二:结合ExternalTaskSensor与额外检查
如果核心需求是等待dag2最新一次运行成功,同时额外确保无运行中实例,可以拆分两个任务:
- 用
ExternalTaskSensor等待dag2的最后一个任务成功 - 用
PythonOperator检查dag2无运行中实例
from airflow.operators.python import PythonOperator from airflow.sensors.external_task import ExternalTaskSensor # 等待dag2的最后一个任务成功 wait_dag2_success = ExternalTaskSensor( task_id="wait_dag2_success", external_dag_id="dag2", external_task_id="dag2_end_task", # 替换为dag2的最后一个任务ID allowed_states=["success"], execution_delta=None, # 等待与当前dag1执行时间匹配的dag2实例,可按需调整 poke_interval=60, mode="reschedule" ) # 检查dag2无运行中实例 def ensure_no_running_dag2(**context): running_runs = DagRun.find(dag_id="dag2", state=DagRunState.RUNNING) if running_runs: raise ValueError("dag2存在运行中实例,终止clean步骤执行") check_no_running = PythonOperator( task_id="check_no_running_dag2", python_callable=ensure_no_running_dag2 ) # 任务流调整为:start >> wait_dag2_success >> check_no_running >> clean >> end
注意事项
- 确保执行dag1的Airflow角色拥有查询dag2运行实例的权限
- 根据实际等待时长调整
poke_interval和timeout参数 mode="reschedule"会在两次检查之间释放worker资源,适合等待时间较长的场景;mode="poke"则持续占用worker资源,适合短时间等待
内容的提问来源于stack exchange,提问作者x89
相关产品推荐
相关产品推荐

