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 64f9cbddc0e Bring back edge worker metric compatibility with Airflow
3.2 (#67328)
64f9cbddc0e is described below
commit 64f9cbddc0ee45db2d5c7872032113a7b725008e
Author: AutomationDev85 <[email protected]>
AuthorDate: Sun Jul 19 23:44:50 2026 +0200
Bring back edge worker metric compatibility with Airflow 3.2 (#67328)
* Bring back edge worker DualStatusManager compability
* Reworked statsd taging for pre airflow 3.3 versions
* Fix mypy
* document metrics compability issues
* fix doc
* Fix docu and added unit test
* Fix imports
* Fix unit test
* Fix flaky unit tests by reset after test execution
---------
Co-authored-by: AutomationDev85 <AutomationDev85>
---
providers/edge3/docs/edge_executor.rst | 27 +++++++++
.../airflow/providers/edge3/models/edge_worker.py | 25 +++++++--
.../tests/unit/edge3/models/test_edge_worker.py | 64 ++++++++++++++++++++++
.../unit/edge3/worker_api/routes/test_worker.py | 17 +++++-
4 files changed, 127 insertions(+), 6 deletions(-)
diff --git a/providers/edge3/docs/edge_executor.rst
b/providers/edge3/docs/edge_executor.rst
index 4452b869589..b29dafbe3af 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 fb800b8b245..ffc6e20383c 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 sqlalchemy.orm import Mapped
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 00000000000..05c6e16a3e3
--- /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 f42d93950f3..939cf5d5f54 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 airflow.providers.edge3.worker_api.routes.worker import (
)
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 @@ class TestWorkerApiRoutes:
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."""