[Airflow]Prevent new tasks from being accidentally run during partial dag reruns

Viewed 167

In airflow, if you add a new task to a dag and then clear its downstream task in an old dag run, airflow will first run the new task you just added.

For example:

A DAG was run on T: A >> B

On T+1 we added a new task C in the middle: A >> C >> B

Now if we don't do anything, and someone clears task B on T, then it will first trigger C, which is not something we want.

Currently we have a manual script to mark a new task as "success" in all historical dag runs, because we don't want the new task to trigger automatically when a downstream task is cleared. I am wondering if there is a less-manual solution to address this concern? It's quite annoying to manually mark success in historical dag runs every time we add a new task.

1 Answers

Since C is in the middle of the workflow I'm not sure there is a simple way to handle this. I would probably create a new DAG with a new start_date. In the old DAG I would add an end_date thus you can backfill both safely if needed.

There might be another solution but i'm not sure this will work well. You can try also to add start_date for task C. Probably you have a star_date defined in default_args where you use it for all your tasks. What you need to do is to overwrite this for task C.

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2020, 4, 1),

}
with DAG(dag_id='stackoverflow',
         default_args=default_args,
         schedule_interval="@daily",
         ) as dag:
    A = DummyOperator(task_id='A')
    B = DummyOperator(task_id='B')
    C = DummyOperator(task_id='C', start_date=datetime(2021, 5, 12))
    
    A >> C >> B

This might be tricky because you define here only the new dependency. This solution would work for sure if the workflow was: A >> B >> C but I never tested it over this use case. In any case my recommendation is to use the first solution of creating a new DAG.

Related