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

如何在Airflow的DAG中实现wait_for_dag2步骤,待DAG2完成后执行clean任务?

实现wait_for_dag2步骤的两种方案

方案一:使用PythonSensor自定义检查逻辑

这种方式可以同时满足「dag2无运行中实例」和「dag2最新实例为成功状态」两个条件,支持周期性检查直到条件满足。

  1. 导入所需模块
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
  1. 编写检查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
  1. 在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最新一次运行成功,同时额外确保无运行中实例,可以拆分两个任务:

  1. 用ExternalTaskSensor等待dag2的最后一个任务成功
  2. 用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 11:40:33