This is an automated email from the ASF dual-hosted git repository.
kaxil 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 370606c1065 Fix Dags staying blocked by `max_active_runs` until the
next parse (#73032)
370606c1065 is described below
commit 370606c1065c7a2f75869d8ecfafdfd9a53a0942
Author: PoAn Yang <[email protected]>
AuthorDate: Tue Oct 6 20:47:50 2026 +0800
Fix Dags staying blocked by `max_active_runs` until the next parse (#73032)
---
airflow-core/src/airflow/jobs/scheduler_job_runner.py | 12 ++----------
airflow-core/tests/unit/jobs/test_scheduler_job.py | 12 ++++++++++--
2 files changed, 12 insertions(+), 12 deletions(-)
diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index 5e07bf3cf70..1f794f41058 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -3156,11 +3156,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
session.flush()
self.log.info("Run %s of %s has timed-out", dag_run.run_id,
dag_run.dag_id)
- if dag_run.state in State.finished_dr_states and dag_run.run_type
in (
- DagRunType.SCHEDULED,
- DagRunType.MANUAL,
- DagRunType.ASSET_TRIGGERED,
- ):
+ if dag_run.state in State.finished_dr_states and dag_run.run_type
!= DagRunType.BACKFILL_JOB:
self._set_exceeds_max_active_runs(dag_model=dag_model,
session=session)
callback_to_execute: DagCallbackRequest | None = None
@@ -3224,11 +3220,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
# TODO[HA]: Rename update_state -> schedule_dag_run, ?? something else?
schedulable_tis, callback_to_run =
dag_run.update_state(session=session, execute_callbacks=False)
- if dag_run.state in State.finished_dr_states and dag_run.run_type in (
- DagRunType.SCHEDULED,
- DagRunType.MANUAL,
- DagRunType.ASSET_TRIGGERED,
- ):
+ if dag_run.state in State.finished_dr_states and dag_run.run_type !=
DagRunType.BACKFILL_JOB:
self._set_exceeds_max_active_runs(dag_model=dag_model,
session=session)
# This will do one query per dag run. We "could" build up a complex
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index eafe68660aa..7211bd68fe4 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -6039,6 +6039,7 @@ class TestSchedulerJob:
)
assert actual == expected
+ @pytest.mark.parametrize("timed_out", [False, True], ids=["finished",
"timed_out"])
@pytest.mark.parametrize(
("run_type", "expected"),
[
@@ -6046,19 +6047,26 @@ class TestSchedulerJob:
(DagRunType.SCHEDULED, True),
(DagRunType.BACKFILL_JOB, False),
(DagRunType.ASSET_TRIGGERED, True),
+ (DagRunType.OPERATOR_TRIGGERED, True),
+ (DagRunType.ASSET_MATERIALIZATION, True),
],
ids=[
DagRunType.MANUAL.name,
DagRunType.SCHEDULED.name,
DagRunType.BACKFILL_JOB.name,
DagRunType.ASSET_TRIGGERED.name,
+ DagRunType.OPERATOR_TRIGGERED.name,
+ DagRunType.ASSET_MATERIALIZATION.name,
],
)
- def test_should_update_dag_next_dagruns_after_run_type(self, run_type,
expected, session, dag_maker):
+ def test_should_update_dag_next_dagruns_after_run_type(
+ self, run_type, expected, timed_out, session, dag_maker
+ ):
"""Test that whether next dag run is updated depends on run type"""
with dag_maker(
schedule="*/1 * * * *",
max_active_runs=3,
+ dagrun_timeout=datetime.timedelta(seconds=60),
):
EmptyOperator(task_id="dummy")
@@ -6066,7 +6074,7 @@ class TestSchedulerJob:
run_id="run",
run_type=run_type,
logical_date=DEFAULT_DATE,
- start_date=timezone.utcnow(),
+ start_date=timezone.utcnow() - datetime.timedelta(days=1) if
timed_out else timezone.utcnow(),
state=State.SUCCESS,
session=session,
)