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

如何实现多层级DAG失败后的自动重跑?含每日调度与历史重试需求

Automated Retry Workflow for Failed DAGs and Tasks

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:

  1. First, create a retry DAG with schedule_interval=None (so it only runs when explicitly triggered).
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 18:32:40