From 39cde663c2d47d604daa71277f7100f033b8d15d Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Fri, 17 Jul 2026 13:56:38 +0200 Subject: [PATCH 1/3] Emit queued and running Dag run counts as metrics --- .../src/airflow/jobs/scheduler_job_runner.py | 17 ++++++++++++----- .../tests/unit/jobs/test_scheduler_job.py | 14 +++++++------- .../observability/metrics/metrics_template.yaml | 12 +++++++++--- 3 files changed, 28 insertions(+), 15 deletions(-) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index f6d6a4787f629..0a931f3c2a3a9 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -1743,7 +1743,7 @@ def _run_scheduler_loop(self) -> None: 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( @@ -3245,10 +3245,17 @@ def _emit_ti_metrics(self, *, session: Session = NEW_SESSION) -> None: 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: + 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}) @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 1f3b6992c3e1c..79c161b258b06 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -9234,8 +9234,8 @@ def test_expired_deadline_locked_by_other_scheduler_is_skipped( 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(self, dag_maker, monkeypatch): + """Test that the dagruns running/queued metrics are emitted correctly.""" with dag_maker("metric_dag") as dag: _ = dag dag_maker.create_dagrun(run_id="run_1", state=DagRunState.RUNNING, logical_date=timezone.utcnow()) @@ -9243,19 +9243,19 @@ def test_emit_running_dags_metric(self, dag_maker, monkeypatch): run_id="run_2", state=DagRunState.RUNNING, logical_date=timezone.utcnow() + timedelta(hours=1) ) - recorded: list[tuple[str, int]] = [] + recorded: list[tuple[str, int, dict]] = [] - def _fake_gauge(metric: str, value: int, *_, **__): - recorded.append((metric, value)) + def _fake_gauge(metric: str, value: int, *_, 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_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, {"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 a3e5b6bdd7555..227470bb8c32f 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -459,10 +459,16 @@ 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, tagged by dag_id." 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, tagged by dag_id." + 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 From debe383657f38af138af1a91efadbdc6b43124c3 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Wed, 5 Aug 2026 08:14:36 +0200 Subject: [PATCH 2/3] Make DagRun metrics configurable --- .../src/airflow/config_templates/config.yml | 13 ++++++ .../src/airflow/jobs/scheduler_job_runner.py | 27 +++++++++--- .../tests/unit/jobs/test_scheduler_job.py | 43 ++++++++++++++++--- .../metrics/metrics_template.yaml | 6 ++- .../observability/metrics/stats.py | 10 +++++ 5 files changed, 84 insertions(+), 15 deletions(-) diff --git a/airflow-core/src/airflow/config_templates/config.yml b/airflow-core/src/airflow/config_templates/config.yml index af11f9fe701d1..aba80eece2569 100644 --- a/airflow-core/src/airflow/config_templates/config.yml +++ b/airflow-core/src/airflow/config_templates/config.yml @@ -2628,6 +2628,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 0a931f3c2a3a9..dd520ed168499 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -3246,16 +3246,29 @@ def _emit_ti_metrics(self, *, session: Session = NEW_SESSION) -> None: @provide_session 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.dag_id, DagRun.state, func.count().label("count")) + select(DagRun.state, func.count().label("count")) .where(DagRun.state.in_([DagRunState.RUNNING, DagRunState.QUEUED])) - .group_by(DagRun.dag_id, DagRun.state) + .group_by(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}) + counts = dict(session.execute(stmt).all()) + 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 79c161b258b06..6f0deff01bdb6 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -9234,28 +9234,59 @@ def test_expired_deadline_locked_by_other_scheduler_is_skipped( mock_handle_miss.assert_not_called() - def test_emit_dag_runs_metric(self, dag_maker, monkeypatch): - """Test that the dagruns running/queued metrics are 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, dict]] = [] + recorded: list[tuple[str, float, dict | None]] = [] - def _fake_gauge(metric: str, value: int, *_, tags=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"}): + 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_dag_runs_metric() - assert recorded == [("scheduler.dagruns.running", 2, {"dag_id": "metric_dag"})] + 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 227470bb8c32f..f5f5f5a63e4ed 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -459,13 +459,15 @@ metrics: name_variables: [] - name: "scheduler.dagruns.running" - description: "Number of DagRuns currently in the ``RUNNING`` state, tagged by dag_id." + 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: "scheduler.dagruns.running.{dag_id}" name_variables: ["dag_id"] - name: "scheduler.dagruns.queued" - description: "Number of DagRuns currently in the ``QUEUED`` state, tagged by dag_id." + 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"] diff --git a/shared/observability/src/airflow_shared/observability/metrics/stats.py b/shared/observability/src/airflow_shared/observability/metrics/stats.py index 5140314922a67..12bc2ce24a1c8 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/stats.py +++ b/shared/observability/src/airflow_shared/observability/metrics/stats.py @@ -142,6 +142,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. From 05554b5730971a20cb866243f8845e4805127c8e Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Wed, 5 Aug 2026 13:25:02 +0200 Subject: [PATCH 3/3] Fix mypy issue --- airflow-core/src/airflow/jobs/scheduler_job_runner.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index dd520ed168499..09e8a55147956 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -3266,7 +3266,9 @@ def _emit_dag_runs_metric(self, *, session: Session = NEW_SESSION) -> None: .where(DagRun.state.in_([DagRunState.RUNNING, DagRunState.QUEUED])) .group_by(DagRun.state) ) - counts = dict(session.execute(stmt).all()) + 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)))