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