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."""

Reply via email to