kjh0623 opened a new issue, #70056: URL: https://github.com/apache/airflow/issues/70056
### Apache Airflow version Observed in production on 3.1.8; the code path is unchanged on 3.3.0 and current main (`b628e46`). ### What happened A `DagRun` created by `TriggerDagRunOperator` sat in `QUEUED` for 8+ hours. The scheduler logged, on every loop: ``` DAG 'YB_DUMMY' not found in serialized_dag table ``` while the DAG itself was healthy: not paused, no import errors, visible in the UI, and new runs of the same DAG worked. Run-scoped API endpoints for the stuck run returned 404 while latest-version endpoints returned 200. The parent DAG's `ExternalTaskSensor` treats only `success`/`failed` as terminal, so it poked the queued run until its own 10h timeout. The chain: 1. The run was created **pinned**: `bundle_version` set, `created_dag_version_id` pointing at the dag version current at creation time. The FK is declared `ondelete="set null"` ([dagrun.py L301-303](https://github.com/apache/airflow/blob/main/airflow-core/src/airflow/models/dagrun.py#L301-L303)). 2. The pinned `dag_version` row was later deleted (the FK's `SET NULL` exists precisely because version rows can be deleted). `created_dag_version_id` became NULL **silently**; `bundle_version` stayed set. 3. `DBDagBag._version_from_dag_run` ([dagbag.py L211-216](https://github.com/apache/airflow/blob/main/airflow-core/src/airflow/models/dagbag.py#L211-L216)) only falls back to the latest version when `bundle_version` is **not** set. For a pinned run it returns `created_dag_version_id` — now NULL — so `get_dag_for_run()` returns `None`. 4. `_start_queued_dagruns` handles that with `log.error(...); continue` ([scheduler_job_runner.py, "DAG '%s' not found in serialized_dag table"](https://github.com/apache/airflow/blob/main/airflow-core/src/airflow/jobs/scheduler_job_runner.py)). The run is neither started nor failed, and is re-selected and re-skipped on every subsequent loop, forever. No timeout applies to a QUEUED run in this state. So a single race between "run created & pinned" and "pinned version row deleted" produces a run that is permanently stuck with no failure signal, no metric, and an error message that points at the wrong cause (the DAG *is* in the serialized_dag table — just not the pinned version). ### What you think should happen instead The scheduler already handles the sibling case for task instances in the same function: when the serialized dag for a SCHEDULED TI cannot be found, the TIs are set to FAILED and the loop moves on. Queued runs deserve the same explicit outcome. Options, in my order of preference: 1. **Fail the run explicitly** when its pinned version cannot be resolved (message naming the pinned `created_dag_version_id`/`bundle_version`), consistent with the TI path. 2. **Fall back to the latest serialized version**, matching unpinned behavior — simpler for users, but arguably violates version-pinning semantics. 3. At minimum, emit a dedicated metric/log so the permanent skip is observable. I'm happy to submit a PR for option 1 (or whichever direction maintainers prefer). ### How to reproduce Two tests against current main — both pass (verified on Linux CI at `b628e46`; the first one "passing" *is* the bug: the run stays QUEUED through repeated scheduler loops). Branch with the tests: https://github.com/kjh0623/airflow/tree/fix/queued-dagrun-stuck-unresolvable-version ```python def test_queued_dagrun_with_unresolvable_pinned_version_is_stuck_forever(self, dag_maker, session): 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() self.job_runner = SchedulerJobRunner(job=Job(), executors=[self.null_exec]) 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.QUEUED # skipped every loop: never started, never failed def test_queued_dagrun_without_bundle_version_falls_back_to_latest(self, dag_maker, session): """Contrast: same NULL created_dag_version_id, but unpinned -> falls back to latest and 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() self.job_runner = SchedulerJobRunner(job=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 ``` On a live deployment: trigger a run into QUEUED (e.g. via `TriggerDagRunOperator` with `wait_for_completion=False` under max_active_runs pressure), then delete its `dag_version` row; the run stays QUEUED indefinitely with the log line above. ### Operating System Kubernetes (observed); reproduction is OS-independent. ### Deployment Official Docker image on Kubernetes, KubernetesExecutor, git-synced bundle (`bundle_version` populated on run creation). ### Anything else Recovery in production required manually clearing the run so a freshly pinned run could be created. SQL to confirm the state: ```sql SELECT run_id, state, bundle_version, created_dag_version_id FROM dag_run WHERE state = 'queued' AND bundle_version IS NOT NULL AND created_dag_version_id IS NULL; ``` ### Willing to submit PR? - [x] Yes I am willing to submit a PR! ### Code of Conduct - [x] I agree to follow this project's [Code of Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
