From cfc97948aa697e69bbf305b98aa28cf220fcc983 Mon Sep 17 00:00:00 2001 From: Prithvi Badiga Date: Thu, 18 Jun 2026 11:25:13 -0700 Subject: [PATCH 1/2] Fix backfill completion race --- .../src/airflow/jobs/scheduler_job_runner.py | 19 ++++- airflow-core/src/airflow/models/backfill.py | 9 +- .../tests/unit/jobs/test_scheduler_job.py | 84 ++++++++++++++++++- 3 files changed, 103 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 93b155f2257ac..035788217dd2a 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -87,7 +87,7 @@ TaskOutletAssetReference, ) from airflow.models.asset_state_store import AssetStateStoreModel -from airflow.models.backfill import Backfill, BackfillDagRun +from airflow.models.backfill import Backfill, BackfillDagRun, BackfillDagRunExceptionReason from airflow.models.callback import Callback, CallbackKey, CallbackType, ExecutorCallback from airflow.models.connection_test import ( ACTIVE_STATES as CONNECTION_TEST_ACTIVE_STATES, @@ -2347,7 +2347,6 @@ def _create_dagruns_for_dags(self, guard: CommitProhibitorGuard, session: Sessio def _mark_backfills_complete(self, *, session: Session = NEW_SESSION) -> None: """Mark completed backfills as completed.""" self.log.debug("checking for completed backfills.") - unfinished_states = (DagRunState.RUNNING, DagRunState.QUEUED) now = timezone.utcnow() # todo: AIP-78 simplify this function to an update statement initializing_cutoff = now - timedelta(minutes=2) @@ -2362,8 +2361,20 @@ def _mark_backfills_complete(self, *, session: Session = NEW_SESSION) -> None: Backfill.created_at < initializing_cutoff, ), ~exists( - select(DagRun.id).where( - and_(DagRun.backfill_id == Backfill.id, DagRun.state.in_(unfinished_states)) + select(BackfillDagRun.id) + .join(DagRun, DagRun.id == BackfillDagRun.dag_run_id, isouter=True) + .where( + BackfillDagRun.backfill_id == Backfill.id, + or_( + and_( + BackfillDagRun.dag_run_id.is_(None), + or_( + BackfillDagRun.exception_reason.is_(None), + BackfillDagRun.exception_reason == BackfillDagRunExceptionReason.IN_FLIGHT, + ), + ), + DagRun.state.in_(State.unfinished_dr_states), + ), ) ), ) diff --git a/airflow-core/src/airflow/models/backfill.py b/airflow-core/src/airflow/models/backfill.py index 52019fa5d984b..eaa8f063830b3 100644 --- a/airflow-core/src/airflow/models/backfill.py +++ b/airflow-core/src/airflow/models/backfill.py @@ -383,7 +383,7 @@ def _create_backfill_dag_run_non_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=None, + dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=non_create_reason, @@ -416,7 +416,7 @@ def _create_backfill_dag_run_non_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=None, + dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=BackfillDagRunExceptionReason.IN_FLIGHT, @@ -460,11 +460,12 @@ def _create_backfill_dag_run_non_partitioned( logical_date=info.logical_date, ) nested.rollback() + latest_dr = session.scalar(_get_latest_dag_run_row_query(dag_id=dag.dag_id, info=info)) session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=None, + dag_run_id=latest_dr.id if latest_dr else None, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=BackfillDagRunExceptionReason.IN_FLIGHT, @@ -492,7 +493,7 @@ def _create_backfill_dag_run_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=None, + dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=non_create_reason, diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index b312f9fd1b0a8..f486fcae11b99 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -67,7 +67,13 @@ AssetPartitionDagRun, PartitionedAssetKeyLog, ) -from airflow.models.backfill import Backfill, BackfillDagRun, ReprocessBehavior, _create_backfill +from airflow.models.backfill import ( + Backfill, + BackfillDagRun, + BackfillDagRunExceptionReason, + ReprocessBehavior, + _create_backfill, +) from airflow.models.callback import Callback, ExecutorCallback from airflow.models.connection_test import ( ConnectionTestKey, @@ -10199,6 +10205,82 @@ def test_mark_backfills_complete_multiple_independent(dag_maker, session): assert b_running.completed_at is None +def test_mark_backfills_complete_waits_for_inflight_association(dag_maker, session): + """Backfill should stay active while an IN_FLIGHT association points at a queued DagRun.""" + clear_db_backfills() + dag_id = "test_backfill_waits_for_inflight_association" + with dag_maker(serialized=True, dag_id=dag_id, schedule="@daily"): + BashOperator(task_id="hi", bash_command="echo hi") + b = Backfill( + dag_id=dag_id, + from_date=pendulum.parse("2021-01-01"), + to_date=pendulum.parse("2021-01-03"), + max_active_runs=10, + dag_run_conf={}, + reprocess_behavior=ReprocessBehavior.NONE, + ) + session.add(b) + session.commit() + backfill_id = b.id + + successful_dr = DagRun( + dag_id=dag_id, + run_id="backfill__2021-01-01T00:00:00+00:00", + run_type=DagRunType.BACKFILL_JOB, + logical_date=pendulum.parse("2021-01-01"), + data_interval=(pendulum.parse("2021-01-01"), pendulum.parse("2021-01-02")), + run_after=pendulum.parse("2021-01-02"), + state=DagRunState.SUCCESS, + backfill_id=backfill_id, + ) + inflight_dr = DagRun( + dag_id=dag_id, + run_id="scheduled__2021-01-02T00:00:00+00:00", + run_type=DagRunType.SCHEDULED, + logical_date=pendulum.parse("2021-01-02"), + data_interval=(pendulum.parse("2021-01-02"), pendulum.parse("2021-01-03")), + run_after=pendulum.parse("2021-01-03"), + state=DagRunState.QUEUED, + ) + session.add_all([successful_dr, inflight_dr]) + session.flush() + session.add_all( + [ + BackfillDagRun( + backfill_id=backfill_id, + dag_run_id=successful_dr.id, + logical_date=pendulum.parse("2021-01-01"), + sort_ordinal=1, + ), + BackfillDagRun( + backfill_id=backfill_id, + dag_run_id=inflight_dr.id, + logical_date=pendulum.parse("2021-01-02"), + exception_reason=BackfillDagRunExceptionReason.IN_FLIGHT, + sort_ordinal=2, + ), + ] + ) + session.commit() + session.expunge_all() + + runner = SchedulerJobRunner( + job=Job(job_type=SchedulerJobRunner.job_type), executors=[MockExecutor(do_update=False)] + ) + runner._mark_backfills_complete() + b = session.get(Backfill, backfill_id) + assert b.completed_at is None + + inflight_dr = session.get(DagRun, inflight_dr.id) + inflight_dr.state = DagRunState.SUCCESS + session.commit() + session.expunge_all() + + runner._mark_backfills_complete() + b = session.get(Backfill, backfill_id) + assert b.completed_at is not None + + class Key1Mapper(CorePartitionMapper): """Partition Mapper that returns only key-1 as downstream key""" From 9e605c49282b08ab6c8488f0f8c62ae14c578b8b Mon Sep 17 00:00:00 2001 From: Prithvi Badiga Date: Thu, 18 Jun 2026 16:11:38 -0700 Subject: [PATCH 2/2] Fix backfill inflight completion tracking --- .../src/airflow/jobs/scheduler_job_runner.py | 26 ++++++++++++++++--- airflow-core/src/airflow/models/backfill.py | 5 ---- 2 files changed, 22 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 035788217dd2a..061073f5583dd 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -2368,12 +2368,30 @@ def _mark_backfills_complete(self, *, session: Session = NEW_SESSION) -> None: or_( and_( BackfillDagRun.dag_run_id.is_(None), - or_( - BackfillDagRun.exception_reason.is_(None), - BackfillDagRun.exception_reason == BackfillDagRunExceptionReason.IN_FLIGHT, - ), + BackfillDagRun.exception_reason.is_(None), ), DagRun.state.in_(State.unfinished_dr_states), + and_( + BackfillDagRun.dag_run_id.is_(None), + BackfillDagRun.exception_reason == BackfillDagRunExceptionReason.IN_FLIGHT, + exists( + select(DagRun.id).where( + DagRun.dag_id == Backfill.dag_id, + DagRun.state.in_(State.unfinished_dr_states), + or_( + and_( + BackfillDagRun.logical_date.is_not(None), + DagRun.logical_date == BackfillDagRun.logical_date, + ), + and_( + BackfillDagRun.logical_date.is_(None), + BackfillDagRun.partition_key.is_not(None), + DagRun.partition_key == BackfillDagRun.partition_key, + ), + ), + ) + ), + ), ), ) ), diff --git a/airflow-core/src/airflow/models/backfill.py b/airflow-core/src/airflow/models/backfill.py index eaa8f063830b3..8453a4a27d60d 100644 --- a/airflow-core/src/airflow/models/backfill.py +++ b/airflow-core/src/airflow/models/backfill.py @@ -383,7 +383,6 @@ def _create_backfill_dag_run_non_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=non_create_reason, @@ -416,7 +415,6 @@ def _create_backfill_dag_run_non_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=BackfillDagRunExceptionReason.IN_FLIGHT, @@ -460,12 +458,10 @@ def _create_backfill_dag_run_non_partitioned( logical_date=info.logical_date, ) nested.rollback() - latest_dr = session.scalar(_get_latest_dag_run_row_query(dag_id=dag.dag_id, info=info)) session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=latest_dr.id if latest_dr else None, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=BackfillDagRunExceptionReason.IN_FLIGHT, @@ -493,7 +489,6 @@ def _create_backfill_dag_run_partitioned( session.add( BackfillDagRun( backfill_id=backfill_id, - dag_run_id=dr.id, logical_date=info.logical_date, partition_key=info.partition_key, exception_reason=non_create_reason,