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 f44634c2154 Detect stale metrics registry entries in the metrics sync
prek hook (#69578)
f44634c2154 is described below
commit f44634c2154e471b92ce294c28e6b8e72a2eb276
Author: PoAn Yang <[email protected]>
AuthorDate: Thu Jul 23 10:47:10 2026 +0900
Detect stale metrics registry entries in the metrics sync prek hook (#69578)
Signed-off-by: PoAn Yang <[email protected]>
---
.pre-commit-config.yaml | 3 +-
.../prek/check_metrics_synced_with_the_registry.py | 162 +++++++++++++++++----
.../test_check_metrics_synced_with_the_registry.py | 129 ++++++++++++++++
.../observability/metrics/metrics_template.yaml | 20 ---
.../observability/metrics/validators.py | 1 -
.../observability/metrics/test_otel_logger.py | 4 +-
6 files changed, 269 insertions(+), 50 deletions(-)
diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml
index 736f3fff51c..adcc9c11c0f 100644
--- a/.pre-commit-config.yaml
+++ b/.pre-commit-config.yaml
@@ -950,9 +950,10 @@ repos:
name: Check that metrics in the codebase are in sync with the metrics
registry YAML file.
entry: ./scripts/ci/prek/check_metrics_synced_with_the_registry.py
language: python
- files: \.py$
+ files:
\.py$|^shared/observability/src/airflow_shared/observability/metrics/metrics_template\.yaml$
exclude:
^(tests/|.*/tests/|airflow-core/src/airflow/observability/stats\.py$|task-sdk/src/airflow/sdk/observability/stats\.py$)
pass_filenames: true
+ require_serial: true
additional_dependencies: ["PyYAML>=6.0", "rich>=13.6.0"]
- id: check-boring-cyborg-configuration
name: Checks for Boring Cyborg configuration consistency
diff --git a/scripts/ci/prek/check_metrics_synced_with_the_registry.py
b/scripts/ci/prek/check_metrics_synced_with_the_registry.py
index 3d42939b6ce..140e6ab9e91 100644
--- a/scripts/ci/prek/check_metrics_synced_with_the_registry.py
+++ b/scripts/ci/prek/check_metrics_synced_with_the_registry.py
@@ -32,6 +32,7 @@ from __future__ import annotations
import argparse
import ast
import re
+import subprocess
import sys
from dataclasses import dataclass
from pathlib import Path
@@ -67,6 +68,16 @@ METRICS_REGISTRY_PATH = (
AIRFLOW_ROOT_PATH /
"shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml"
)
+EXCLUDED_TESTS_PATTERN = re.compile(r"(^|/)tests/")
+
+# Registry metrics emitted through a variable that the AST scan cannot resolve
back to a
+# string literal, e.g. names built by BaseExecutor._get_metric_name.
+INDIRECTLY_EMITTED_METRICS = {
+ "executor.open_slots",
+ "executor.queued_tasks",
+ "executor.running_tasks",
+}
+
def load_metrics_registry_yaml() -> dict[str, dict[str, Any]]:
"""Load the metrics registry YAML and return a dict keyed by metric
name."""
@@ -92,6 +103,12 @@ def normalize_metric_name(registry_metric_name: str) -> str:
_PREFIX_MATCHED = "__prefix_matched__"
+def find_prefix_matched_registry_entries(metric_name: str, metrics_registry:
dict[str, dict]) -> list[str]:
+ """Return the registry entry names whose name matches the static prefix of
a dynamic metric name."""
+ base = metric_name.split("{")[0].rstrip(".")
+ return [name for name in metrics_registry if name == base or
name.startswith(base + ".")]
+
+
def find_registry_match(metric_name: str, metrics_registry: dict[str, dict])
-> str | None:
"""Return the registry entry name that best matches the given metric name,
or None."""
if metric_name in metrics_registry:
@@ -109,16 +126,13 @@ def find_registry_match(metric_name: str,
metrics_registry: dict[str, dict]) ->
return registry_metric_name
# Dynamic metric name.
- if "{" in metric_name:
- base = metric_name.split("{")[0].rstrip(".")
- for registry_metric_name in metrics_registry:
- if registry_metric_name == base or
registry_metric_name.startswith(base + "."):
- # Metric prefix matches the prefix of a dynamic registry entry.
- # If the static part before the first variable, matches an
exact registry entry name,
- # or a dotted-prefix of one, then τηε name is considered
covered and
- # _PREFIX_MATCHED is returned. The type check must be skipped
because
- # the resulting metric name with all variables expanded,
cannot be determined.
- return _PREFIX_MATCHED
+ if "{" in metric_name and
find_prefix_matched_registry_entries(metric_name, metrics_registry):
+ # Metric prefix matches the prefix of a dynamic registry entry.
+ # If the static part before the first variable, matches an exact
registry entry name,
+ # or a dotted-prefix of one, then the name is considered covered and
+ # _PREFIX_MATCHED is returned. The type check must be skipped because
+ # the resulting metric name with all variables expanded, cannot be
determined.
+ return _PREFIX_MATCHED
# All checks for matching failed.
return None
@@ -177,6 +191,22 @@ def extract_metric_name_from_ast_node(name_node: ast.expr)
-> str | None:
return None
+def extract_metric_names_from_ast_node(name_node: ast.expr) -> list[str]:
+ """Resolve a metric name AST node to all names it can produce.
+
+ A conditional expression like ``"a.success" if ok else "a.failed"`` yields
both names.
+ Returns an empty list when the expression is too dynamic to resolve.
+ """
+ if isinstance(name_node, ast.IfExp):
+ return [
+ name
+ for branch in (name_node.body, name_node.orelse)
+ for name in extract_metric_names_from_ast_node(branch)
+ ]
+ metric_name = extract_metric_name_from_ast_node(name_node)
+ return [metric_name] if metric_name is not None else []
+
+
@dataclass
class MetricCall:
file_path: str
@@ -232,8 +262,16 @@ def scan_file_for_metrics(file_path: Path) ->
list[MetricCall]:
"""Return all Stats metric calls found in the provided file_path."""
try:
source = file_path.read_text(encoding="utf-8")
+ except (OSError, UnicodeDecodeError):
+ return []
+
+ # Cheap text pre-filter to avoid parsing the vast majority of files that
emit no metrics.
+ if "Stats." not in source and "stats." not in source:
+ return []
+
+ try:
tree = ast.parse(source, filename=str(file_path))
- except (OSError, UnicodeDecodeError, SyntaxError):
+ except SyntaxError:
return []
metrics_found: list[MetricCall] = []
@@ -251,22 +289,22 @@ def scan_file_for_metrics(file_path: Path) ->
list[MetricCall]:
if name_node is None:
continue
- metric_name = extract_metric_name_from_ast_node(name_node)
- if metric_name is None:
+ if not (metric_names :=
extract_metric_names_from_ast_node(name_node)):
# Metric name is unresolvable. Probably has too many variables.
continue
stats_obj = get_stats_obj_name(node.func.value)
- metrics_found.append(
- MetricCall(
- file_path=str(file_path),
- line_num=node.lineno,
- metric_name=metric_name,
- method=method,
- stats_obj=stats_obj or "",
- is_dynamic="{" in metric_name,
+ for metric_name in metric_names:
+ metrics_found.append(
+ MetricCall(
+ file_path=str(file_path),
+ line_num=node.lineno,
+ metric_name=metric_name,
+ method=method,
+ stats_obj=stats_obj or "",
+ is_dynamic="{" in metric_name,
+ )
)
- )
return metrics_found
@@ -318,6 +356,43 @@ def scan_file_for_direct_stats_imports(file_path: Path) ->
list[DirectStatsImpor
return violations
+def list_repository_python_files() -> list[Path]:
+ """Return all git-tracked, non-test Python files in the repository."""
+ result = subprocess.run(
+ ["git", "ls-files", "*.py"],
+ cwd=AIRFLOW_ROOT_PATH,
+ capture_output=True,
+ text=True,
+ check=True,
+ )
+ return [
+ AIRFLOW_ROOT_PATH / line
+ for line in result.stdout.splitlines()
+ if not EXCLUDED_TESTS_PATTERN.search(line)
+ ]
+
+
+def compute_unused_registry_entries(
+ code_metric_names: set[str], metrics_registry: dict[str, dict]
+) -> list[str]:
+ """Return the registry entry names that no code metric name matches."""
+ used_entries = set(INDIRECTLY_EMITTED_METRICS)
+ for metric_name in code_metric_names:
+ registry_metric_name = find_registry_match(metric_name,
metrics_registry)
+ if registry_metric_name is None:
+ continue
+ if registry_metric_name is _PREFIX_MATCHED:
+
used_entries.update(find_prefix_matched_registry_entries(metric_name,
metrics_registry))
+ else:
+ used_entries.add(registry_metric_name)
+ return sorted(set(metrics_registry) - used_entries)
+
+
+def find_stale_indirectly_emitted_metrics(metrics_registry: dict[str, dict])
-> list[str]:
+ """Return the ``INDIRECTLY_EMITTED_METRICS`` names that no longer exist in
the registry."""
+ return sorted(INDIRECTLY_EMITTED_METRICS - set(metrics_registry))
+
+
def main() -> None:
parser = argparse.ArgumentParser(
description="Check that metrics in the codebase are in sync with the
metrics registry YAML file."
@@ -377,12 +452,24 @@ def main() -> None:
if mismatched:
metrics_with_type_mismatch[name] = mismatched
- # There is no point in checking whether the metrics exist in the YAML but
not in the code,
- # because the script is comparing the entire YAML against certain files at
a time.
- # For that to work, the script would have to run against all project files
EVERY TIME.
+ # Violation 3: the metric exists in the registry but nowhere in the code.
The hook only
+ # receives the changed files, so this check scans all git-tracked Python
files itself.
+ all_code_metric_names = {
+ call.metric_name
+ for repo_file_path in list_repository_python_files()
+ for call in scan_file_for_metrics(repo_file_path)
+ }
+ unused_registry_entries =
compute_unused_registry_entries(all_code_metric_names, metrics_registry)
+
+ # Violation 4: this script exempts a metric from the check above, but the
registry no longer has it.
+ stale_indirectly_emitted_metrics =
find_stale_indirectly_emitted_metrics(metrics_registry)
total_violations = (
- len(metrics_not_in_registry) + len(metrics_with_type_mismatch) +
len(direct_stats_imports)
+ len(metrics_not_in_registry)
+ + len(metrics_with_type_mismatch)
+ + len(direct_stats_imports)
+ + len(stale_indirectly_emitted_metrics)
+ + len(unused_registry_entries)
)
if total_violations:
@@ -436,6 +523,29 @@ def main() -> None:
)
console.print()
+ if unused_registry_entries:
+ console.print(
+ f" [red]-> {len(unused_registry_entries)} metric(s) found
in the registry YAML but not emitted anywhere in the code:[/red]"
+ )
+ for metric_name in unused_registry_entries:
+ console.print(f" [green]{metric_name}[/green]")
+ console.print(
+ " [yellow]Remove them from the registry, or if they are
emitted through a variable "
+ "the scan cannot resolve, add them to
INDIRECTLY_EMITTED_METRICS in this script.[/yellow]"
+ )
+ console.print()
+
+ if stale_indirectly_emitted_metrics:
+ console.print(
+ f" [red]-> {len(stale_indirectly_emitted_metrics)}
metric(s) in INDIRECTLY_EMITTED_METRICS that no longer exist in the registry
YAML:[/red]"
+ )
+ for metric_name in stale_indirectly_emitted_metrics:
+ console.print(f" [green]{metric_name}[/green]")
+ console.print(
+ " [yellow]Update INDIRECTLY_EMITTED_METRICS in this script
to match the registry.[/yellow]"
+ )
+ console.print()
+
sys.exit(1)
diff --git
a/scripts/tests/ci/prek/test_check_metrics_synced_with_the_registry.py
b/scripts/tests/ci/prek/test_check_metrics_synced_with_the_registry.py
index fdd91fb4dce..e9b5e2ebe29 100644
--- a/scripts/tests/ci/prek/test_check_metrics_synced_with_the_registry.py
+++ b/scripts/tests/ci/prek/test_check_metrics_synced_with_the_registry.py
@@ -19,14 +19,20 @@ from __future__ import annotations
import ast
import textwrap
from pathlib import Path
+from unittest import mock
import pytest
+from ci.prek import check_metrics_synced_with_the_registry
from ci.prek.check_metrics_synced_with_the_registry import (
_PREFIX_MATCHED,
_except_handler_catches_expected_error,
_is_stats_module_path,
+ compute_unused_registry_entries,
extract_metric_name_from_ast_node,
+ extract_metric_names_from_ast_node,
+ find_prefix_matched_registry_entries,
find_registry_match,
+ find_stale_indirectly_emitted_metrics,
get_stats_obj_name,
normalize_metric_name,
scan_file_for_direct_stats_imports,
@@ -183,6 +189,121 @@ def test_extract_metric_name_from_ast_node(code: str,
expected_result):
assert extract_metric_name_from_ast_node(node) == expected_result
[email protected](
+ "code, expected_result",
+ [
+ pytest.param('"scheduler_heartbeat"', ["scheduler_heartbeat"],
id="static_string_single_name"),
+ pytest.param(
+ '"connection_test.success" if success else
"connection_test.failed"',
+ ["connection_test.success", "connection_test.failed"],
+ id="conditional_expression_both_branches",
+ ),
+ pytest.param(
+ '"connection_test.success" if success else get_name()',
+ ["connection_test.success"],
+ id="conditional_expression_unresolvable_branch_dropped",
+ ),
+ pytest.param(
+ '"a" if x else ("b" if y else "c")',
+ ["a", "b", "c"],
+ id="nested_conditional_expression",
+ ),
+ pytest.param("some_variable", [],
id="unresolvable_name_returns_empty_list"),
+ ],
+)
+def test_extract_metric_names_from_ast_node(code: str, expected_result):
+ node = ast.parse(code, mode="eval").body
+ assert extract_metric_names_from_ast_node(node) == expected_result
+
+
[email protected](
+ "metric_name, expected_result",
+ [
+ pytest.param(
+ "ti.{state}",
+ ["ti.scheduled", "ti.queued", "ti.start.{dag_id}.{task_id}"],
+ id="base_prefix_matches_multiple_entries",
+ ),
+ pytest.param("dagrun.duration.{state}", ["dagrun.duration.success"],
id="dotted_base_prefix"),
+ pytest.param("non.existent.{var}", [],
id="no_prefix_match_returns_empty_list"),
+ ],
+)
+def test_find_prefix_matched_registry_entries(metric_name, expected_result):
+ assert find_prefix_matched_registry_entries(metric_name, METRICS_REGISTRY)
== expected_result
+
+
+# 'executor.open_slots' is in INDIRECTLY_EMITTED_METRICS, so it is never
reported as unused.
[email protected](
+ "code_metric_names, expected_unused",
+ [
+ pytest.param(
+ set(),
+ [
+ "dagrun.duration.success",
+ "pool.open_slots",
+ "scheduler.heartbeat",
+ "task.duration",
+ "ti.queued",
+ "ti.scheduled",
+ "ti.start.{dag_id}.{task_id}",
+ ],
+ id="no_code_metrics_reports_all_but_indirectly_emitted",
+ ),
+ pytest.param(
+ {"scheduler.heartbeat", "dag.{x}.{y}.duration", "unknown.metric"},
+ [
+ "dagrun.duration.success",
+ "pool.open_slots",
+ "ti.queued",
+ "ti.scheduled",
+ "ti.start.{dag_id}.{task_id}",
+ ],
+ id="exact_and_legacy_matches_mark_entries_used",
+ ),
+ pytest.param(
+ {"ti.{state}"},
+ ["dagrun.duration.success", "pool.open_slots",
"scheduler.heartbeat", "task.duration"],
+ id="prefix_match_marks_all_prefix_entries_used",
+ ),
+ pytest.param(
+ {"pool.open_slots.{my_pool}"},
+ [
+ "dagrun.duration.success",
+ "scheduler.heartbeat",
+ "task.duration",
+ "ti.queued",
+ "ti.scheduled",
+ "ti.start.{dag_id}.{task_id}",
+ ],
+ id="legacy_name_structure_match_marks_entry_used",
+ ),
+ ],
+)
+def test_compute_unused_registry_entries(code_metric_names, expected_unused):
+ assert compute_unused_registry_entries(code_metric_names,
METRICS_REGISTRY) == expected_unused
+
+
[email protected](
+ "indirectly_emitted_metrics, expected_stale",
+ [
+ pytest.param({"executor.open_slots"}, [],
id="entry_present_in_registry"),
+ pytest.param(set(), [], id="empty_allowlist"),
+ pytest.param(
+ {"executor.open_slots", "executor.renamed_away",
"executor.deleted"},
+ ["executor.deleted", "executor.renamed_away"],
+ id="renamed_and_deleted_entries_reported",
+ ),
+ ],
+)
+def test_find_stale_indirectly_emitted_metrics(indirectly_emitted_metrics,
expected_stale):
+ with mock.patch.object(
+ check_metrics_synced_with_the_registry,
+ "INDIRECTLY_EMITTED_METRICS",
+ indirectly_emitted_metrics,
+ ):
+ assert find_stale_indirectly_emitted_metrics(METRICS_REGISTRY) ==
expected_stale
+
+
@pytest.fixture
def code_to_py_file(tmp_path):
"""Write python source code to a tmp file and return its path."""
@@ -240,6 +361,14 @@ def code_to_py_file(tmp_path):
[{"stats_obj": "stats"}],
id="self_stats_attribute",
),
+ pytest.param(
+ 'stats.incr("connection_test.success" if success else
"connection_test.failed")',
+ [
+ {"metric_name": "connection_test.success", "method": "incr"},
+ {"metric_name": "connection_test.failed", "method": "incr"},
+ ],
+ id="conditional_expression_yields_call_per_branch",
+ ),
pytest.param('metrics.incr("triggerer_heartbeat")', [],
id="unknown_stats_object_ignored"),
pytest.param("Stats.incr(get_metric_name())", [],
id="unresolvable_metric_name_skipped"),
pytest.param("def foo(:\n pass\n", [],
id="syntax_error_returns_empty"),
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 c0d43641559..a561db2d7ae 100644
---
a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
+++
b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml
@@ -39,14 +39,6 @@ metrics:
legacy_name: "-"
name_variables: ["job_name"]
- - name: "local_task_job.task_exit"
- description: "Number of ``LocalTaskJob`` terminations with a
``{return_code}`` while
- running a task ``{task_id}`` of a Dag ``{dag_id}``.
- Metric with job_id, dag_id, task_id and return_code tagging."
- type: "counter"
- legacy_name:
"local_task_job.task_exit.{job_id}.{dag_id}.{task_id}.{return_code}"
- name_variables: ["job_id", "dag_id", "task_id", "return_code"]
-
- name: "operator_failures"
description: "Operator ``{operator_name}`` failures."
type: "counter"
@@ -162,12 +154,6 @@ metrics:
legacy_name: "-"
name_variables: []
- - name: "dag_file_processor_timeouts"
- description: "(DEPRECATED) same behavior as
``dag_processing.processor_timeouts``"
- type: "counter"
- legacy_name: "-"
- name_variables: []
-
- name: "scheduler.tasks.killed_externally"
description: "Number of tasks killed externally. Metric with dag_id and
task_id tagging."
type: "counter"
@@ -433,12 +419,6 @@ metrics:
legacy_name: "dag_processing.last_run.seconds_ago.{file_name}"
name_variables: ["file_name"]
- - name: "dag_processing.last_num_of_db_queries.{dag_file}"
- description: "Number of queries to Airflow database during parsing per
``{dag_file}``"
- type: "gauge"
- legacy_name: "-"
- name_variables: ["dag_file"]
-
- name: "scheduler.tasks.starving"
description: "Number of tasks that cannot be scheduled because of no open
slot in pool"
type: "gauge"
diff --git
a/shared/observability/src/airflow_shared/observability/metrics/validators.py
b/shared/observability/src/airflow_shared/observability/metrics/validators.py
index 3d9d3d3d3e5..cab81c7d98b 100644
---
a/shared/observability/src/airflow_shared/observability/metrics/validators.py
+++
b/shared/observability/src/airflow_shared/observability/metrics/validators.py
@@ -56,7 +56,6 @@ BACK_COMPAT_METRIC_NAME_PATTERNS: set[str] = {
r"^(?P<job_name>.*)_start$",
r"^(?P<job_name>.*)_end$",
r"^(?P<job_name>.*)_heartbeat_failure$",
-
r"^local_task_job.task_exit\.(?P<job_id>.*)\.(?P<dag_id>.*)\.(?P<task_id>.*)\.(?P<return_code>.*)$",
r"^operator_failures_(?P<operator_name>.*)$",
r"^operator_successes_(?P<operator_name>.*)$",
r"^ti.start.(?P<dag_id>.*)\.(?P<task_id>.*)$",
diff --git
a/shared/observability/tests/observability/metrics/test_otel_logger.py
b/shared/observability/tests/observability/metrics/test_otel_logger.py
index 2b5c71193c3..f1e90e53003 100644
--- a/shared/observability/tests/observability/metrics/test_otel_logger.py
+++ b/shared/observability/tests/observability/metrics/test_otel_logger.py
@@ -77,9 +77,9 @@ class TestOtelMetrics:
assert not _is_up_down_counter("this_is_not_a_udc")
def test_exemption_list_has_not_grown(self):
- assert len(BACK_COMPAT_METRIC_NAMES) <= 26, (
+ assert len(BACK_COMPAT_METRIC_NAMES) <= 25, (
"This test exists solely to ensure that nobody is adding names to
the exemption list. "
- "There are 26 names which are potentially too long for OTel and
that number should "
+ "There are 25 names which are potentially too long for OTel and
that number should "
"only ever go down as these names are deprecated. If this test is
failing, please "
"adjust your new stat's name; do not add as exemption without a
very good reason."
)