如何实现多层级DAG失败后的自动重跑?含每日调度与历史重试需求
Absolutely! This is totally doable (especially if you're using Apache Airflow, given the DAG terminology you mentioned). Let's break down how to build this automated retry workflow step by step:
1. Auto-Retry Failed Tasks the Next Day
You have two solid approaches here, depending on how much control you need:
Approach 1: Task-Level Retry with Scheduled Delay
If your main DAGs run daily (schedule_interval='@daily'), you can leverage Airflow's built-in retry mechanism to configure each task to retry exactly once with a 24-hour delay:
from datetime import timedelta from airflow.operators.python import PythonOperator def my_task_function(**context): # Your core task logic here pass my_task = PythonOperator( task_id='my_task', python_callable=my_task_function, retries=1, # Only retry once retry_delay=timedelta(days=1), # Wait 24 hours before retrying execution_timeout=timedelta(hours=2), # Adjust based on your task's typical runtime dag=dag )
This retry will be tied to the original DAG Run, so it will execute automatically the next day when the retry window opens.
Approach 2: Trigger a Dedicated Retry DAG on Failure
For more flexibility (like separating retry logic from your main workflow), use an on_failure_callback to trigger a standalone retry DAG when a task fails. This retry DAG can be set to run the next day:
- First, create a retry DAG with
schedule_interval=None(so it only runs when explicitly triggered). - Add a callback function to your main tasks:
from airflow.api.client.local_client import Client from datetime import datetime, timedelta def trigger_retry_dag(context): client = Client(None, None) # Calculate the next day's execution date next_day = context['execution_date'] + timedelta(days=1) # Trigger the retry DAG with context about the failed task client.trigger_dag( dag_id='retry_failed_task_dag', run_id=f"retry_{context['task_instance_key_str']}", execution_date=next_day, conf={ 'failed_task_id': context['task_instance'].task_id, 'original_dag_id': context['task_instance'].dag_id } ) # Attach the callback to your task my_task = PythonOperator( task_id='my_task', python_callable=my_task_function, on_failure_callback=trigger_retry_dag, dag=dag )
The retry DAG can then use the passed conf values to rerun the specific failed task.
2. Auto-Rerun All Failed DAGs from the Last 10 Days
Create a dedicated cleanup DAG that runs daily to scan for and retrigger failed DAG runs from the past 10 days:
Step 1: Build the Cleanup DAG
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.models import DagRun from airflow.utils.state import State from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'retries': 0 } def retry_failed_dags(**context): # Define the 10-day lookback window ten_days_ago = datetime.now() - timedelta(days=10) # Query all failed DAG runs in the window failed_runs = DagRun.query.filter( DagRun.state == State.FAILED, DagRun.execution_date >= ten_days_ago ).all() # Retrigger each failed run (avoid duplicates) for run in failed_runs: new_run_id = f"retry_{run.run_id}" # Check if this retry run already exists to prevent loops existing_run = DagRun.query.filter(DagRun.run_id == new_run_id).first() if not existing_run: dag = context['dag'].get_dag(run.dag_id) dag.create_dagrun( run_id=new_run_id, execution_date=run.execution_date, state=State.RUNNING, conf=run.conf ) print(f"Retriggered failed DAG Run: {run.dag_id} - {run.execution_date}") with DAG( 'retry_failed_dags_daily', default_args=default_args, description='Daily DAG to retry failed DAG runs from the last 10 days', schedule_interval='@daily', catchup=False ) as dag: retry_task = PythonOperator( task_id='retry_failed_dags', python_callable=retry_failed_dags, provide_context=True )
Step 2: Configure Permissions
Make sure the Airflow service account running this cleanup DAG has permission to trigger other DAGs. You can set this up in the Airflow UI under Security > Permissions to grant the necessary access.
Key Considerations
- Resource Management: Retrying 10 days of failed DAGs can consume significant resources. Schedule the cleanup DAG to run during off-peak hours, and set concurrency limits on your DAGs if needed.
- Logging: Add detailed logging to your retry logic so you can track which tasks/DAGs were retried and their outcomes.
- Loop Prevention: Always check if a retry run already exists before triggering it to avoid infinite retry loops.
内容的提问来源于stack exchange,提问作者user2221179

