This is an automated email from the ASF dual-hosted git repository.
Lee-W pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 4d356a64d57 Skip timeout callback creation for Dags without
on_failure_callback (#70065)
4d356a64d57 is described below
commit 4d356a64d5702829c123756c292b6bc4281da1cf
Author: PoAn Yang <[email protected]>
AuthorDate: Wed Sep 16 08:24:07 2026 +0900
Skip timeout callback creation for Dags without on_failure_callback (#70065)
---
.../src/airflow/jobs/scheduler_job_runner.py | 41 ++++++++++++----------
airflow-core/tests/unit/jobs/test_scheduler_job.py | 32 +++++++++++++++++
2 files changed, 54 insertions(+), 19 deletions(-)
diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index fe1e09a9ee5..90aefb30be9 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -3040,26 +3040,29 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
):
self._set_exceeds_max_active_runs(dag_model=dag_model,
session=session)
- dag_run_reloaded = session.scalar(
- select(DagRun)
- .where(DagRun.id == dag_run.id)
- .options(
-
selectinload(DagRun.consumed_asset_events).selectinload(AssetEvent.asset),
-
selectinload(DagRun.consumed_asset_events).selectinload(AssetEvent.source_aliases),
+ callback_to_execute: DagCallbackRequest | None = None
+ if dag.has_on_failure_callback:
+ # Only load the asset events when a callback will actually be
produced.
+ dag_run_reloaded = session.scalar(
+ select(DagRun)
+ .where(DagRun.id == dag_run.id)
+ .options(
+
selectinload(DagRun.consumed_asset_events).selectinload(AssetEvent.asset),
+
selectinload(DagRun.consumed_asset_events).selectinload(AssetEvent.source_aliases),
+ )
+ )
+ if dag_run_reloaded is None:
+ # This should never happen since we just had the dag_run
+ self.log.error("DagRun %s was deleted unexpectedly",
dag_run.id)
+ return None
+ dag_run = dag_run_reloaded
+ callback_to_execute = dag_run.produce_dag_callback(
+ dag=dag,
+ success=False,
+ relevant_ti=last_unfinished_ti,
+ reason="timed_out",
+ execute=False,
)
- )
- if dag_run_reloaded is None:
- # This should never happen since we just had the dag_run
- self.log.error("DagRun %s was deleted unexpectedly",
dag_run.id)
- return None
- dag_run = dag_run_reloaded
- callback_to_execute = dag_run.produce_dag_callback(
- dag=dag,
- success=False,
- relevant_ti=last_unfinished_ti,
- reason="timed_out",
- execute=False,
- )
# Team name should be added before listeners are called in
notify_dagrun_state_changed()
self._stamp_team_names([dag_run], session)
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index ca91ef5f76f..4d1a540feda 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -1152,6 +1152,7 @@ class TestSchedulerJob:
schedule=[asset1],
fileloc="/test_path1/",
dagrun_timeout=timedelta(minutes=1),
+ on_failure_callback=lambda ctx: None,
):
EmptyOperator(task_id="dummy_task")
@@ -4190,6 +4191,7 @@ class TestSchedulerJob:
start_date=DEFAULT_DATE,
max_active_runs=1,
dagrun_timeout=datetime.timedelta(seconds=60),
+ on_failure_callback=lambda ctx: None,
) as dag:
EmptyOperator(task_id="dummy")
@@ -4254,6 +4256,7 @@ class TestSchedulerJob:
with dag_maker(
dag_id="test_scheduler_fail_dagrun_timeout",
dagrun_timeout=datetime.timedelta(seconds=60),
+ on_failure_callback=lambda ctx: None,
session=session,
):
EmptyOperator(task_id="dummy")
@@ -4315,6 +4318,35 @@ class TestSchedulerJob:
session.rollback()
session.close()
+ @mock.patch.object(DagRun, "produce_dag_callback", autospec=True)
+ def test_dagrun_timeout_without_on_failure_callback_produces_no_callback(
+ self, mock_produce_dag_callback, dag_maker
+ ):
+ """A timed-out run of a Dag without on_failure_callback must not build
a callback request."""
+ session = settings.Session()
+ with dag_maker(
+ dag_id="test_scheduler_dagrun_timeout_no_callback",
+ dagrun_timeout=datetime.timedelta(seconds=60),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy")
+
+ dr = dag_maker.create_dagrun(start_date=timezone.utcnow() -
datetime.timedelta(days=1))
+
+ scheduler_job = Job()
+ self.job_runner = SchedulerJobRunner(job=scheduler_job)
+
+ callback = self.job_runner._schedule_dag_run(dr, session)
+ session.flush()
+
+ session.refresh(dr)
+ assert dr.state == State.FAILED
+ assert callback is None
+ mock_produce_dag_callback.assert_not_called()
+
+ session.rollback()
+ session.close()
+
def test_dagrun_timeout_fails_run_and_update_next_dagrun(self, dag_maker):
"""
Test that dagrun timeout fails run and update the next dagrun