Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions providers/edge3/docs/edge_executor.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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."""
Expand Down
64 changes: 64 additions & 0 deletions providers/edge3/tests/unit/edge3/models/test_edge_worker.py
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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."""
Expand Down