diff --git a/providers/edge3/docs/edge_executor.rst b/providers/edge3/docs/edge_executor.rst index 4452b86958966..b29dafbe3afbd 100644 --- a/providers/edge3/docs/edge_executor.rst +++ b/providers/edge3/docs/edge_executor.rst @@ -217,3 +217,30 @@ Current Limitations Edge Executor - Multi-team isolation is logical only — all teams share a single authentication secret. A worker administrator could change the team name and access another team's jobs. See :ref:`edge_executor:multi_team` for details and planned improvements. + + +Metrics Export Compatibility +----------------------------- + +The Edge Worker integrates with Airflow's metrics system to export runtime metrics. Compatibility +between Edge provider versions and Airflow versions varies due to changes in the metrics initialization +pipeline. The table below documents known compatibility issues and workarounds: + +.. list-table:: + :header-rows: 1 + + * - Provider version + - Airflow 3.3 + - Airflow 3.2 + - Airflow 3.1 + * - >= 3.6.0 + - Working + - Broken: webserver missing a `stats.initialize(...)` call -- DualStatsManager removed, missing legacy metric names + - Broken: metric tags + * - <= 3.5.0 + - Working + - Requires manual patch of ``metrics_template.yaml`` to match previous export schema + - Working + +**For Airflow 3.2 users:** If upgrading to Edge provider >= 3.6.0 breaks metrics export, either +(1) upgrade Airflow to 3.3+, or (2) downgrade to Edge provider <= 3.5.0 with the workaround above. diff --git a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py index fb800b8b2457b..ffc6e20383cde 100644 --- a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py +++ b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py @@ -28,6 +28,7 @@ from airflow.providers.common.compat.sdk import AirflowException, Stats, timezone from airflow.providers.common.compat.sqlalchemy.orm import mapped_column from airflow.providers.edge3.models.edge_base import Base +from airflow.providers.edge3.version_compat import AIRFLOW_V_3_3_PLUS from airflow.utils.helpers import prune_dict from airflow.utils.log.logging_mixin import LoggingMixin from airflow.utils.providers_configuration_loader import providers_configuration_loaded @@ -181,12 +182,11 @@ def set_metrics( "free_concurrency", } metric_tags = prune_dict({"worker_name": worker_name, "team_name": team_name}) + status = sysinfo.get("status", logging.NOTSET) + if not isinstance(status, (int, float)): + status = logging.NOTSET - Stats.gauge( - "edge_worker.status", - sysinfo.get("status", logging.NOTSET), # type: ignore - tags=metric_tags, - ) + Stats.gauge("edge_worker.status", status, tags=metric_tags) Stats.gauge("edge_worker.connected", int(connected), tags=metric_tags) Stats.gauge("edge_worker.maintenance", int(maintenance), tags=metric_tags) Stats.gauge("edge_worker.jobs_active", jobs_active, tags=metric_tags) @@ -203,6 +203,21 @@ def set_metrics( if isinstance(value, (int, float)): Stats.gauge(f"edge_worker.{key}", value, tags=metric_tags) + if not AIRFLOW_V_3_3_PLUS: + # Airflow < 3.3: export legacy per-worker metrics (no auto-tag expansion). + Stats.gauge(f"edge_worker.status.{worker_name}", int(status)) + Stats.gauge(f"edge_worker.connected.{worker_name}", int(connected)) + Stats.gauge(f"edge_worker.maintenance.{worker_name}", int(maintenance)) + Stats.gauge(f"edge_worker.jobs_active.{worker_name}", jobs_active) + Stats.gauge(f"edge_worker.concurrency.{worker_name}", concurrency) + Stats.gauge(f"edge_worker.free_concurrency.{worker_name}", free_concurrency) + Stats.gauge(f"edge_worker.num_queues.{worker_name}", len(queues)) + + for key in additional_keys: + value = sysinfo.get(key) + if isinstance(value, (int, float)): + Stats.gauge(f"edge_worker.{key}.{worker_name}", value) + def reset_metrics(worker_name: str, team_name: str | None = None) -> None: """Reset metrics of worker.""" diff --git a/providers/edge3/tests/unit/edge3/models/test_edge_worker.py b/providers/edge3/tests/unit/edge3/models/test_edge_worker.py new file mode 100644 index 0000000000000..05c6e16a3e33b --- /dev/null +++ b/providers/edge3/tests/unit/edge3/models/test_edge_worker.py @@ -0,0 +1,64 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from __future__ import annotations + +from unittest import mock + +from airflow.providers.common.compat.sdk import Stats +from airflow.providers.edge3.models.edge_worker import EdgeWorkerState, set_metrics + +from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS + +stats_reference = f"{Stats.__module__}.Stats" + + +def test_set_metrics(): + worker_name = "test_worker1" + if AIRFLOW_V_3_3_PLUS: + with mock.patch("airflow.sdk._shared.observability.metrics.stats._get_backend") as mock_get_backend: + mock_backend = mock.MagicMock() + mock_get_backend.return_value = mock_backend + + set_metrics( + worker_name=worker_name, + state=EdgeWorkerState.IDLE, + jobs_active=0, + concurrency=1, + free_concurrency=1, + queues=None, + sysinfo={"status": 1}, + ) + + metric_names = [call.args[0] for call in mock_backend.gauge.call_args_list] + else: + with mock.patch(f"{stats_reference}.gauge") as mock_gauge: + set_metrics( + worker_name=worker_name, + state=EdgeWorkerState.IDLE, + jobs_active=0, + concurrency=1, + free_concurrency=1, + queues=None, + sysinfo={"status": 1}, + ) + + metric_names = [call.args[0] for call in mock_gauge.call_args_list] + + assert "edge_worker.status" in metric_names + + legacy_metric_name = f"edge_worker.status.{worker_name}" + assert legacy_metric_name in metric_names diff --git a/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py b/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py index f42d93950f3b7..939cf5d5f54ff 100644 --- a/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py +++ b/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py @@ -45,6 +45,7 @@ ) from tests_common.test_utils.config import conf_vars +from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS if TYPE_CHECKING: from sqlalchemy.orm import Session @@ -423,7 +424,21 @@ def test_set_state_metrics_team_name_tags( tags={**expected_worker_tags, "queues": ",".join(queues)}, ) mock_stats_gauge.assert_any_call("edge_worker.disk_usage", 42.5, tags=expected_worker_tags) - assert mock_stats_gauge.call_count == 8 + if AIRFLOW_V_3_3_PLUS: + assert mock_stats_gauge.call_count == 8 + else: + mock_stats_gauge.assert_any_call( + "edge_worker.status.test2_worker", + self.MOCK_SYSINFO["status"], + ) + mock_stats_gauge.assert_any_call("edge_worker.connected.test2_worker", 1) + mock_stats_gauge.assert_any_call("edge_worker.maintenance.test2_worker", 0) + mock_stats_gauge.assert_any_call("edge_worker.jobs_active.test2_worker", 1) + mock_stats_gauge.assert_any_call("edge_worker.concurrency.test2_worker", 8) + mock_stats_gauge.assert_any_call("edge_worker.free_concurrency.test2_worker", 8) + mock_stats_gauge.assert_any_call("edge_worker.num_queues.test2_worker", len(queues)) + mock_stats_gauge.assert_any_call("edge_worker.disk_usage.test2_worker", 42.5) + assert mock_stats_gauge.call_count == 16 def test_set_state_returns_concurrency(self, session: Session, cli_worker: EdgeWorker): """set_state includes the DB-stored concurrency override in its response."""