如何让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.
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"}'
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

