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

如何让Airflow在当前DAGRun完成后自动禁用DAG?

Great questions! Let's break down both of your requirements with practical, actionable solutions using Airflow's CLI, API, and built-in workflow tools.

1. Directly Fail the Latest DAGRun

You can achieve this via either the Airflow CLI or REST API, depending on what fits your workflow better.

Using the CLI

Airflow's dags set-run-state command lets you directly update the state of a specific DAGRun. To target the latest running instance, use the --latest flag:

airflow dags set-run-state --dag-id your_target_dag_id --latest --state failed

This will immediately mark the most recent DAGRun as failed, regardless of its current execution state—perfect for your use case where you don't care about preserving its progress.

Using the REST API

If you prefer programmatic control, you'll need two steps: first fetch the latest DAGRun ID, then update its state to failed.

Step 1: Get the latest DAGRun ID

Send a GET request to retrieve the most recent run (sorted by execution date descending):

# Replace placeholders with your Airflow details
LATEST_RUN_ID=$(curl -X GET "http://your-airflow-webserver:8080/api/v1/dags/your_target_dag_id/dagRuns?limit=1&order_by=-execution_date" \
  -H "Authorization: Bearer your_auth_token" \
  | jq -r '.dag_runs[0].dag_run_id')

(Note: Skip the Authorization header if your Airflow instance doesn't use authentication.)

Step 2: Set the DAGRun state to failed

Use the retrieved ID to send a POST request updating the state:

curl -X POST "http://your-airflow-webserver:8080/api/v1/dags/your_target_dag_id/dagRuns/$LATEST_RUN_ID/set_state" \
  -H "Authorization: Bearer your_auth_token" \
  -H "Content-Type: application/json" \
  -d '{"state": "failed"}'
2. Auto-Disable DAG After Current DAGRun Completes

To ensure the DAG is paused only after the current run finishes (whether it succeeds or fails), you have two reliable approaches:

Method 1: Add a Final "Disable DAG" Task

Add a dedicated task at the end of your DAG that runs regardless of previous task outcomes (using TriggerRule.ALL_DONE). This task will pause the DAG once executed.

Option A: BashOperator (uses Airflow CLI)

from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime

with DAG(
    dag_id="your_target_dag_id",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,  # Adjust to your schedule
) as dag:
    # Your existing tasks go here
    process_data = ...

    # Final task to pause the DAG
    disable_dag = BashOperator(
        task_id="disable_dag_after_run",
        bash_command="airflow dags pause your_target_dag_id",
        trigger_rule=TriggerRule.ALL_DONE,  # Runs even if prior tasks fail
    )

    process_data >> disable_dag

Option B: PythonOperator (direct Airflow API call)

For a more integrated approach without relying on the CLI, use Airflow's internal models to pause the DAG:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.trigger_rule import TriggerRule
from airflow.models.dag import DagModel
from datetime import datetime

def pause_dag():
    dag_model = DagModel.get_dagmodel("your_target_dag_id")
    dag_model.is_paused = True
    dag_model.update()

with DAG(
    dag_id="your_target_dag_id",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
) as dag:
    # Your existing tasks
    process_data = ...

    disable_dag = PythonOperator(
        task_id="disable_dag_after_run",
        python_callable=pause_dag,
        trigger_rule=TriggerRule.ALL_DONE,
    )

    process_data >> disable_dag

Method 2: Use DAG-Level Callbacks

If you don't want to add an extra task, leverage DAG-level success/failure callbacks to trigger the disable action whenever the run completes:

from airflow import DAG
from airflow.models.dag import DagModel
from datetime import datetime

def disable_dag_on_completion(context):
    dag_id = context["dag"].dag_id
    dag_model = DagModel.get_dagmodel(dag_id)
    dag_model.is_paused = True
    dag_model.update()

with DAG(
    dag_id="your_target_dag_id",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    on_success_callback=disable_dag_on_completion,
    on_failure_callback=disable_dag_on_completion,
) as dag:
    # Your tasks here

This callback will run whenever the DAGRun finishes (whether it succeeds or fails), automatically pausing the DAG afterward.


内容的提问来源于stack exchange,提问作者yee379

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:21:18