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

Airflow中如何基于全局变量控制DAG任务执行或触发DAG?

解决方案

一、在clean任务前等待STATUS变量变为True

可以用Airflow的PythonSensor实现等待逻辑,把它插入到start和clean任务之间,循环检查全局变量状态,直到满足条件再继续执行后续任务。

示例代码:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.sensors.python import PythonSensor
from airflow.models import Variable
from datetime import datetime

def start_task():
    print("Start task executed")

def clean_task():
    print("Clean task executed")

def end_task():
    print("End task executed")

def check_status():
    # Airflow变量默认存储为字符串,需做类型转换
    status = Variable.get("STATUS", default_var="false").lower()
    return status == "true"

with DAG(
    dag_id="dag1",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # 不需要固定调度可设为None
    catchup=False
) as dag:
    start = PythonOperator(
        task_id="start",
        python_callable=start_task
    )

    check_status_sensor = PythonSensor(
        task_id="check_status_sensor",
        python_callable=check_status,
        poke_interval=30,  # 每30秒检查一次
        timeout=86400,  # 最长等待1天,超时则任务失败
        mode="poke"  # 等待久的话推荐用reschedule模式,释放worker资源
    )

    clean = PythonOperator(
        task_id="clean",
        python_callable=clean_task
    )

    end = PythonOperator(
        task_id="end",
        python_callable=end_task
    )

    start >> check_status_sensor >> clean >> end

说明:

  • poke_interval:可根据实际需求调整变量检查间隔
  • mode="reschedule":适合长等待场景,传感器会暂时释放worker资源,到下一次检查时间再重新调度

二、当STATUS变量设为True时触发dag1

如果不想让DAG处于等待状态,可通过以下两种方式实现变量触发:

1. 利用Airflow Listener(Airflow 2.2+)

自定义Listener监听变量更新事件,当STATUS被设为true时自动触发dag1。将以下代码放在Airflow的plugins目录下即可生效:

from airflow.listeners.listener import Listener
from airflow.models import Variable
from airflow.utils.session import provide_session
from airflow.api.common.experimental.trigger_dag import trigger_dag

class VariableTriggerListener(Listener):
    @provide_session
    def on_variable_updated(self, session, variable, previous_variable):
        if variable.key == "STATUS" and variable.value.lower() == "true":
            # 触发目标DAG
            trigger_dag(dag_id="dag1", run_id=None, conf=None, session=session)

# 注册Listener
VariableTriggerListener()

2. 手动/脚本触发

在设置STATUS变量为true的同时,调用Airflow的CLI或API触发dag1:

  • CLI方式:
airflow variables set STATUS true && airflow dags trigger dag1
  • 若通过UI设置变量,可配合监控脚本定时检查STATUS状态,当变为true时触发DAG。

内容的提问来源于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:55:20