From 1d76037ff44a1abe91d6bbad2ce20888843be088 Mon Sep 17 00:00:00 2001 From: kjh0623 <8412070+kjh0623@users.noreply.github.com> Date: Mon, 20 Jul 2026 12:34:37 +0900 Subject: [PATCH 1/3] Fail queued DagRuns whose pinned dag version can no longer be resolved A DagRun created pinned to a dag version (bundle_version set, created_dag_version_id populated) becomes unresolvable if that dag_version row is later deleted: the FK's ondelete="set null" clears created_dag_version_id while bundle_version stays set, so the pinned resolver returns None (it only falls back to the latest version for unpinned runs). _start_queued_dagruns then could not resolve the serialized DAG and skipped the run with "DAG '...' not found in serialized_dag table" on every loop - the run was neither started nor failed and stayed QUEUED indefinitely, with no failure signal and a misleading log (the DAG is in serialized_dag; the pinned version is not). Fail such a run explicitly, mirroring how a SCHEDULED task instance whose serialized DAG cannot be found is already failed in the same loop. Add a regression test that the pinned-unresolvable run reaches FAILED instead of being skipped forever, plus a contrast test that an unpinned run with the same NULL created_dag_version_id still falls back to the latest version and starts. Signed-off-by: kjh0623 <8412070+kjh0623@users.noreply.github.com> --- airflow-core/newsfragments/70056.bugfix.rst | 3 ++ .../src/airflow/jobs/scheduler_job_runner.py | 17 +++++- .../tests/unit/jobs/test_scheduler_job.py | 54 +++++++++++++++++++ 3 files changed, 73 insertions(+), 1 deletion(-) create mode 100644 airflow-core/newsfragments/70056.bugfix.rst diff --git a/airflow-core/newsfragments/70056.bugfix.rst b/airflow-core/newsfragments/70056.bugfix.rst new file mode 100644 index 0000000000000..81bf613aeb8e5 --- /dev/null +++ b/airflow-core/newsfragments/70056.bugfix.rst @@ -0,0 +1,3 @@ +Fail queued DagRuns whose pinned dag version can no longer be resolved, instead of skipping them forever. + +When a ``DagRun`` is created pinned to a dag version (``bundle_version`` set and ``created_dag_version_id`` populated) and that ``dag_version`` row is later deleted, the foreign key's ``ondelete="set null"`` clears ``created_dag_version_id`` while ``bundle_version`` stays set. The pinned resolver then returns ``None`` (it only falls back to the latest version for unpinned runs), so the scheduler could not resolve the serialized DAG for the queued run and skipped it with a ``DAG '...' not found in serialized_dag table`` log line on every loop -- the run was never started nor failed and stayed ``QUEUED`` indefinitely with no failure signal. The scheduler now fails such a run explicitly, consistent with how a scheduled task instance whose serialized DAG cannot be found is already handled. diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index c2f3ee05d730e..1771f7aa61b30 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -2768,7 +2768,22 @@ def _update_state(dag: SerializedDAG, dag_run: DagRun): backfill_id = dag_run.backfill_id dag = dag_run.dag = cached_get_dag(dag_run) if not dag: - self.log.error("DAG '%s' not found in serialized_dag table", dag_run.dag_id) + # The serialized DAG for this run cannot be resolved. This happens when the run is + # pinned to a dag version (bundle_version set, created_dag_version_id populated) whose + # dag_version row was later deleted: the FK's ondelete="set null" clears + # created_dag_version_id and the pinned resolver does not fall back to the latest + # version. Without an explicit outcome the run is re-selected and skipped on every + # loop forever, with no failure signal. Fail it explicitly, mirroring how a SCHEDULED + # task instance whose serialized DAG cannot be found is failed elsewhere in the loop. + self.log.error( + "DAG '%s' for queued run %s could not be resolved " + "(created_dag_version_id=%s, bundle_version=%s); marking the run as failed", + dag_id, + run_id, + dag_run.created_dag_version_id, + dag_run.bundle_version, + ) + dag_run.set_state(DagRunState.FAILED) continue active_runs = active_runs_of_dags[(dag_id, backfill_id)] if backfill_id is not None: diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index 68f4ec0915092..ce4930f21de3c 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -6007,6 +6007,60 @@ def test_start_dagruns(self, mock_get_backend, dag_maker, session): assert get_last_dagrun(dag.dag_id, session).creating_job_id == scheduler_job.id + def test_queued_dagrun_with_unresolvable_pinned_version_is_failed(self, dag_maker, session): + """ + A queued run whose pinned dag version can no longer be resolved is failed + explicitly instead of being skipped on every scheduler loop. + + Reproduces the state after the pinned dag_version row is deleted: the FK + ondelete="set null" clears created_dag_version_id while bundle_version + remains set, so the pinned resolver returns None and never falls back to + the latest version. The scheduler now marks the run FAILED rather than + re-selecting and skipping it forever. + """ + with dag_maker(dag_id="test_stuck_pinned_run"): + EmptyOperator(task_id="mytask") + dr = dag_maker.create_dagrun(run_type=DagRunType.MANUAL, state=State.QUEUED) + + # The run was created pinned to a bundle version... + dr.bundle_version = "0123456789abcdef" + # ...and the pinned dag_version row has since been deleted -> FK SET NULL + dr.created_dag_version_id = None + session.merge(dr) + session.flush() + + scheduler_job = Job() + self.job_runner = SchedulerJobRunner(job=scheduler_job, executors=[self.null_exec]) + + # Even across repeated loops the run reaches a terminal state instead of being skipped. + for _ in range(3): + self.job_runner._start_queued_dagruns(session) + session.flush() + + dr = session.scalars(select(DagRun).where(DagRun.dag_id == "test_stuck_pinned_run")).one() + assert dr.state == State.FAILED + + def test_queued_dagrun_without_bundle_version_falls_back_to_latest(self, dag_maker, session): + """ + Contrast case: same NULL created_dag_version_id, but without bundle_version + the resolver falls back to the latest serialized version and the run starts. + """ + with dag_maker(dag_id="test_unpinned_run_falls_back"): + EmptyOperator(task_id="mytask") + dr = dag_maker.create_dagrun(run_type=DagRunType.MANUAL, state=State.QUEUED) + dr.bundle_version = None + dr.created_dag_version_id = None + session.merge(dr) + session.flush() + + scheduler_job = Job() + self.job_runner = SchedulerJobRunner(job=scheduler_job, executors=[self.null_exec]) + self.job_runner._start_queued_dagruns(session) + session.flush() + + dr = session.scalars(select(DagRun).where(DagRun.dag_id == "test_unpinned_run_falls_back")).one() + assert dr.state == State.RUNNING + def test_extra_operator_links_not_loaded_in_scheduler_loop(self, dag_maker): """ Test that Operator links are not loaded inside the Scheduling Loop (that does not include From ac0d7ece971077b6aa7197f50d248dd2dbfcddf6 Mon Sep 17 00:00:00 2001 From: kjh0623 <8412070+kjh0623@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:27:27 +0900 Subject: [PATCH 2/3] Make the newsfragment a single line check-newsfragments-are-valid requires non-significant newsfragments to be a single line; the detailed explanation belongs in the PR description, not the newsfragment. Signed-off-by: kjh0623 <8412070+kjh0623@users.noreply.github.com> --- airflow-core/newsfragments/70056.bugfix.rst | 2 -- 1 file changed, 2 deletions(-) diff --git a/airflow-core/newsfragments/70056.bugfix.rst b/airflow-core/newsfragments/70056.bugfix.rst index 81bf613aeb8e5..a7383671da696 100644 --- a/airflow-core/newsfragments/70056.bugfix.rst +++ b/airflow-core/newsfragments/70056.bugfix.rst @@ -1,3 +1 @@ Fail queued DagRuns whose pinned dag version can no longer be resolved, instead of skipping them forever. - -When a ``DagRun`` is created pinned to a dag version (``bundle_version`` set and ``created_dag_version_id`` populated) and that ``dag_version`` row is later deleted, the foreign key's ``ondelete="set null"`` clears ``created_dag_version_id`` while ``bundle_version`` stays set. The pinned resolver then returns ``None`` (it only falls back to the latest version for unpinned runs), so the scheduler could not resolve the serialized DAG for the queued run and skipped it with a ``DAG '...' not found in serialized_dag table`` log line on every loop -- the run was never started nor failed and stayed ``QUEUED`` indefinitely with no failure signal. The scheduler now fails such a run explicitly, consistent with how a scheduled task instance whose serialized DAG cannot be found is already handled. From b0d056bb43e9024d6a3e3cfa6ff079d3ac797d82 Mon Sep 17 00:00:00 2001 From: kjh0623 <8412070+kjh0623@users.noreply.github.com> Date: Tue, 21 Jul 2026 22:33:50 +0900 Subject: [PATCH 3/3] Name the newsfragment after the PR number check-newsfragment-pr-number requires the file to be named after the PR, not the issue it fixes: 70056.bugfix.rst -> 70113.bugfix.rst. Signed-off-by: kjh0623 <8412070+kjh0623@users.noreply.github.com> --- airflow-core/newsfragments/{70056.bugfix.rst => 70113.bugfix.rst} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename airflow-core/newsfragments/{70056.bugfix.rst => 70113.bugfix.rst} (100%) diff --git a/airflow-core/newsfragments/70056.bugfix.rst b/airflow-core/newsfragments/70113.bugfix.rst similarity index 100% rename from airflow-core/newsfragments/70056.bugfix.rst rename to airflow-core/newsfragments/70113.bugfix.rst