This is an automated email from the ASF dual-hosted git repository.

potiuk 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 88043d0782e Improve Dag run metrics (#70013)
88043d0782e is described below

commit 88043d0782e328ad2dbecd76917b033bd02cfdd0
Author: AutomationDev85 <[email protected]>
AuthorDate: Tue Sep 8 23:34:50 2026 +0200

    Improve Dag run metrics (#70013)
    
    * Emit queued and running Dag run counts as metrics
    
    * Make DagRun metrics configurable
    
    * Fix mypy issue
    
    * Add unit test and newfragment
    
    ---------
    
    Co-authored-by: AutomationDev85 <AutomationDev85>
---
 airflow-core/newsfragments/70013.feature.rst       |  1 +
 .../src/airflow/config_templates/config.yml        | 13 ++++++
 .../src/airflow/jobs/scheduler_job_runner.py       | 32 ++++++++++++---
 airflow-core/tests/unit/jobs/test_scheduler_job.py | 47 ++++++++++++++++++----
 .../observability/metrics/metrics_template.yaml    | 14 +++++--
 .../airflow_shared/observability/metrics/stats.py  | 10 +++++
 .../tests/observability/metrics/test_stats.py      |  8 +++-
 7 files changed, 108 insertions(+), 17 deletions(-)

diff --git a/airflow-core/newsfragments/70013.feature.rst 
b/airflow-core/newsfragments/70013.feature.rst
new file mode 100644
index 00000000000..428e22d49ea
--- /dev/null
+++ b/airflow-core/newsfragments/70013.feature.rst
@@ -0,0 +1 @@
+Added a new ``scheduler.dagruns.queued`` metric alongside the existing 
``scheduler.dagruns.running`` metric, so operators can track queued DagRun 
backlog (e.g. Dags stuck behind ``max_active_runs`` limits or pool exhaustion) 
as well as running load. Both metrics are emitted as a single untagged 
aggregate value across all Dags by default, matching the previous behaviour of 
``scheduler.dagruns.running``. Set ``[scheduler] dagrun_metrics_per_dag_id`` to 
``True`` to instead emit one gauge pe [...]
diff --git a/airflow-core/src/airflow/config_templates/config.yml 
b/airflow-core/src/airflow/config_templates/config.yml
index 0a0dc060a0a..3601a94759c 100644
--- a/airflow-core/src/airflow/config_templates/config.yml
+++ b/airflow-core/src/airflow/config_templates/config.yml
@@ -2725,6 +2725,19 @@ scheduler:
       type: float
       example: ~
       default: "30.0"
+    dagrun_metrics_per_dag_id:
+      description: |
+        If true, the ``scheduler.dagruns.running`` and 
``scheduler.dagruns.queued`` metrics are
+        emitted once per ``dag_id`` (tagged by ``dag_id``) instead of as a 
single aggregate value
+        across all Dags. This gives per-Dag visibility into running/queued 
backlog, at the cost of
+        emitting one metric series per Dag with an active DagRun on every
+        ``[scheduler] dagrun_metrics_interval``. On deployments with a large 
number of Dags, enabling
+        this can significantly increase the number of metric series sent to 
StatsD/OpenTelemetry, so
+        it is disabled by default.
+      version_added: 3.4.0
+      type: boolean
+      example: ~
+      default: "False"
     scheduler_health_check_threshold:
       description: |
         If the last scheduler heartbeat happened more than ``[scheduler] 
scheduler_health_check_threshold``
diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py 
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index 127f385f0e0..a082334f2a1 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -1845,7 +1845,7 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
 
             timers.call_regular_interval(
                 conf.getfloat("scheduler", "dagrun_metrics_interval", 
fallback=30.0),
-                self._emit_running_dags_metric,
+                self._emit_dag_runs_metric,
             )
 
         timers.call_regular_interval(
@@ -3357,10 +3357,32 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
             self.previous_ti_metrics[state] = ti_metrics
 
     @provide_session
-    def _emit_running_dags_metric(self, *, session: Session = NEW_SESSION) -> 
None:
-        stmt = select(func.count()).select_from(DagRun).where(DagRun.state == 
DagRunState.RUNNING)
-        running_dags = float(session.scalar(stmt) or 0)
-        stats.gauge("scheduler.dagruns.running", running_dags)
+    def _emit_dag_runs_metric(self, *, session: Session = NEW_SESSION) -> None:
+        if conf.getboolean("scheduler", "dagrun_metrics_per_dag_id"):
+            stmt = (
+                select(DagRun.dag_id, DagRun.state, 
func.count().label("count"))
+                .where(DagRun.state.in_([DagRunState.RUNNING, 
DagRunState.QUEUED]))
+                .group_by(DagRun.dag_id, DagRun.state)
+            )
+            for dag_id, state, count in session.execute(stmt).all():
+                metric_name = (
+                    "scheduler.dagruns.running"
+                    if state == DagRunState.RUNNING
+                    else "scheduler.dagruns.queued"
+                )
+                stats.gauge(metric_name, float(count), tags={"dag_id": dag_id})
+            return
+
+        stmt = (
+            select(DagRun.state, func.count().label("count"))
+            .where(DagRun.state.in_([DagRunState.RUNNING, DagRunState.QUEUED]))
+            .group_by(DagRun.state)
+        )
+        counts: dict[DagRunState, int] = {}
+        for state, count in session.execute(stmt):
+            counts[state] = int(count)
+        stats.gauge("scheduler.dagruns.running", 
float(counts.get(DagRunState.RUNNING, 0)))
+        stats.gauge("scheduler.dagruns.queued", 
float(counts.get(DagRunState.QUEUED, 0)))
 
     @provide_session
     def _emit_pool_metrics(self, *, session: Session = NEW_SESSION) -> None:
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py 
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index a7b332e763f..d9ce14a35e4 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -10178,28 +10178,59 @@ class TestSchedulerJob:
 
             mock_handle_miss.assert_not_called()
 
-    def test_emit_running_dags_metric(self, dag_maker, monkeypatch):
-        """Test that the running_dags metric is emitted correctly."""
+    def test_emit_dag_runs_metric_aggregate_by_default(self, dag_maker, 
monkeypatch):
+        """Test that the dagruns running/queued metrics are emitted as 
untagged aggregates by default."""
         with dag_maker("metric_dag") as dag:
             _ = dag
         dag_maker.create_dagrun(run_id="run_1", state=DagRunState.RUNNING, 
logical_date=timezone.utcnow())
         dag_maker.create_dagrun(
             run_id="run_2", state=DagRunState.RUNNING, 
logical_date=timezone.utcnow() + timedelta(hours=1)
         )
+        dag_maker.create_dagrun(
+            run_id="run_3", state=DagRunState.QUEUED, 
logical_date=timezone.utcnow() + timedelta(hours=2)
+        )
 
-        recorded: list[tuple[str, int]] = []
+        recorded: list[tuple[str, float, dict | None]] = []
 
-        def _fake_gauge(metric: str, value: int, *_, **__):
-            recorded.append((metric, value))
+        def _fake_gauge(metric: str, value: float, *_, tags=None, **__):
+            recorded.append((metric, value, tags))
 
         
monkeypatch.setattr("airflow._shared.observability.metrics.stats.gauge", 
_fake_gauge, raising=True)
 
-        with conf_vars({("metrics", "statsd_on"): "True"}):
+        with conf_vars(
+            {("metrics", "statsd_on"): "True", ("scheduler", 
"dagrun_metrics_per_dag_id"): "False"}
+        ):
+            scheduler_job = Job()
+            self.job_runner = SchedulerJobRunner(scheduler_job)
+            self.job_runner._emit_dag_runs_metric()
+
+        assert ("scheduler.dagruns.running", 2.0, None) in recorded
+        assert ("scheduler.dagruns.queued", 1.0, None) in recorded
+
+    def test_emit_dag_runs_metric_per_dag_id_when_enabled(self, dag_maker, 
monkeypatch):
+        """Test that the dagruns running/queued metrics are tagged by dag_id 
when opted in."""
+        with dag_maker("metric_dag") as dag:
+            _ = dag
+        dag_maker.create_dagrun(run_id="run_1", state=DagRunState.RUNNING, 
logical_date=timezone.utcnow())
+        dag_maker.create_dagrun(
+            run_id="run_2", state=DagRunState.RUNNING, 
logical_date=timezone.utcnow() + timedelta(hours=1)
+        )
+
+        recorded: list[tuple[str, float, dict | None]] = []
+
+        def _fake_gauge(metric: str, value: float, *_, tags=None, **__):
+            recorded.append((metric, value, tags))
+
+        
monkeypatch.setattr("airflow._shared.observability.metrics.stats.gauge", 
_fake_gauge, raising=True)
+
+        with conf_vars(
+            {("metrics", "statsd_on"): "True", ("scheduler", 
"dagrun_metrics_per_dag_id"): "True"}
+        ):
             scheduler_job = Job()
             self.job_runner = SchedulerJobRunner(scheduler_job)
-            self.job_runner._emit_running_dags_metric()
+            self.job_runner._emit_dag_runs_metric()
 
-        assert recorded == [("scheduler.dagruns.running", 2)]
+        assert recorded == [("scheduler.dagruns.running", 2.0, {"dag_id": 
"metric_dag"})]
 
     # Multi-team scheduling tests
     def test_multi_team_get_team_names_for_dag_ids_success(self, dag_maker, 
session):
diff --git 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
index b4a39ea4fb5..6cbcf3dc304 100644
--- 
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
+++ 
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
@@ -495,10 +495,18 @@ metrics:
     name_variables: []
 
   - name: "scheduler.dagruns.running"
-    description: "Number of DAGs whose latest DagRun is currently in the 
``RUNNING`` state"
+    description: "Number of DagRuns currently in the ``RUNNING`` state. 
Emitted as a single aggregate
+    value by default; tagged by dag_id instead when ``[scheduler] 
dagrun_metrics_per_dag_id`` is enabled."
     type: "gauge"
-    legacy_name: "-"
-    name_variables: []
+    legacy_name: "scheduler.dagruns.running.{dag_id}"
+    name_variables: ["dag_id"]
+
+  - name: "scheduler.dagruns.queued"
+    description: "Number of DagRuns currently in the ``QUEUED`` state. Emitted 
as a single aggregate
+    value by default; tagged by dag_id instead when ``[scheduler] 
dagrun_metrics_per_dag_id`` is enabled."
+    type: "gauge"
+    legacy_name: "scheduler.dagruns.queued.{dag_id}"
+    name_variables: ["dag_id"]
 
   - name: "executor.open_slots"
     description: "Number of open slots on executor. Legacy metric only emitted
diff --git 
a/shared/observability/src/airflow_shared/observability/metrics/stats.py 
b/shared/observability/src/airflow_shared/observability/metrics/stats.py
index 4904a18d0b6..e083f340dc1 100644
--- a/shared/observability/src/airflow_shared/observability/metrics/stats.py
+++ b/shared/observability/src/airflow_shared/observability/metrics/stats.py
@@ -161,6 +161,16 @@ def _get_legacy_stat_name_and_tags(
         return _none
 
     required_vars = stat_from_registry.get("name_variables", [])
+
+    # tags=None means this call opted out of tagging (e.g. an untagged 
aggregate
+    # behind a config flag), so skip the legacy name instead of raising.
+    # Example: ``scheduler.dagruns.running`` uses legacy
+    # ``scheduler.dagruns.running.{dag_id}``; when emitted as an aggregate with
+    # ``tags=None``, we do not try to format ``{dag_id}``. An empty dict still
+    # raises below since that means tags were expected but missing.
+    if required_vars and tags is None:
+        return _none
+
     provided_vars = set(tags.keys()) if tags else set()
     missing_vars = set(required_vars) - provided_vars
     # If there are specified variables in the YAML file that haven't been 
provided in the tags param.
diff --git a/shared/observability/tests/observability/metrics/test_stats.py 
b/shared/observability/tests/observability/metrics/test_stats.py
index ee3e9dc8a81..c2a4a931021 100644
--- a/shared/observability/tests/observability/metrics/test_stats.py
+++ b/shared/observability/tests/observability/metrics/test_stats.py
@@ -644,12 +644,18 @@ class TestStatsHelpers:
                 (None, {}),
                 id="missing_metric_returns_none",
             ),
+            pytest.param(
+                "operator_failures",
+                None,
+                (None, {}),
+                id="tags_none_returns_no_legacy_name",
+            ),
         ],
     )
     def test_get_legacy_stat_name_and_tags(
         self,
         stat: str,
-        tags: dict,
+        tags: dict | None,
         expected_result: tuple[str | None, dict],
     ):
         from airflow_shared.observability.metrics.stats import 
_get_legacy_stat_name_and_tags

Reply via email to