Skip to content

Commit 256ae43

Browse files
committed
replace try-except import pattern with AIRFLOW_V_3_2_PLUS checks
1 parent 94c06b1 commit 256ae43

8 files changed

Lines changed: 32 additions & 38 deletions

File tree

providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
from airflow.executors.base_executor import BaseExecutor
3030
from airflow.models.taskinstance import TaskInstance
3131
from airflow.providers.common.compat.sdk import Stats, timezone
32+
from airflow.providers.common.compat.version_compat import AIRFLOW_V_3_2_PLUS
3233
from airflow.providers.edge3.models.db import EdgeDBManager, check_db_manager_config
3334
from airflow.providers.edge3.models.edge_job import EdgeJobModel
3435
from airflow.providers.edge3.models.edge_logs import EdgeLogsModel
@@ -37,11 +38,6 @@
3738
from airflow.utils.session import NEW_SESSION, provide_session
3839
from airflow.utils.state import TaskInstanceState
3940

40-
try:
41-
from airflow.sdk.observability.stats import DualStatsManager
42-
except ImportError:
43-
DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat
44-
4541
if TYPE_CHECKING:
4642
from sqlalchemy.orm import Session
4743

@@ -197,7 +193,9 @@ def _update_orphaned_jobs(self, session: Session) -> bool:
197193
"queue": job.queue,
198194
"state": str(TaskInstanceState.FAILED),
199195
}
200-
if DualStatsManager is not None:
196+
if AIRFLOW_V_3_2_PLUS:
197+
from airflow.sdk.observability.stats import DualStatsManager
198+
201199
DualStatsManager.incr("edge_worker.ti.finish", tags={}, legacy_name_tags=tags)
202200
else:
203201
Stats.incr(

providers/edge3/src/airflow/providers/edge3/models/edge_worker.py

Lines changed: 5 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,9 @@
2727
from sqlalchemy.orm import Mapped
2828

2929
from airflow.providers.common.compat.sdk import AirflowException, Stats, timezone
30-
from airflow.providers.edge3.models.edge_base import Base
31-
32-
try:
33-
from airflow.sdk.observability.stats import DualStatsManager
34-
except ImportError:
35-
DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat
3630
from airflow.providers.common.compat.sqlalchemy.orm import mapped_column
31+
from airflow.providers.common.compat.version_compat import AIRFLOW_V_3_2_PLUS
32+
from airflow.providers.edge3.models.edge_base import Base
3733
from airflow.utils.log.logging_mixin import LoggingMixin
3834
from airflow.utils.providers_configuration_loader import providers_configuration_loaded
3935
from airflow.utils.session import NEW_SESSION, provide_session
@@ -178,7 +174,9 @@ def set_metrics(
178174
EdgeWorkerState.OFFLINE_MAINTENANCE,
179175
)
180176

181-
if DualStatsManager is not None:
177+
if AIRFLOW_V_3_2_PLUS:
178+
from airflow.sdk.observability.stats import DualStatsManager
179+
182180
DualStatsManager.gauge(
183181
"edge_worker.connected",
184182
int(connected),

providers/edge3/src/airflow/providers/edge3/worker_api/routes/jobs.py

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,7 @@
2727
from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
2828
from airflow.executors.workloads import ExecuteTask
2929
from airflow.providers.common.compat.sdk import Stats, timezone
30-
31-
try:
32-
from airflow.sdk.observability.stats import DualStatsManager
33-
except ImportError:
34-
DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat
30+
from airflow.providers.common.compat.version_compat import AIRFLOW_V_3_2_PLUS
3531
from airflow.providers.edge3.models.edge_job import EdgeJobModel
3632
from airflow.providers.edge3.worker_api.auth import jwt_token_authorization_rest
3733
from airflow.providers.edge3.worker_api.datamodels import (
@@ -41,6 +37,9 @@
4137
)
4238
from airflow.utils.state import TaskInstanceState
4339

40+
if AIRFLOW_V_3_2_PLUS:
41+
from airflow.sdk.observability.stats import DualStatsManager
42+
4443
jobs_router = AirflowRouter(tags=["Jobs"], prefix="/jobs")
4544

4645

@@ -91,7 +90,7 @@ def fetch(
9190
session.commit()
9291
# Edge worker does not backport emitted Airflow metrics, so export some metrics
9392
tags = {"dag_id": job.dag_id, "task_id": job.task_id, "queue": job.queue}
94-
if DualStatsManager is not None:
93+
if AIRFLOW_V_3_2_PLUS:
9594
DualStatsManager.incr("edge_worker.ti.start", tags={}, legacy_name_tags=tags)
9695
else:
9796
Stats.incr(f"edge_worker.ti.start.{job.queue}.{job.dag_id}.{job.task_id}", tags=tags)
@@ -148,7 +147,7 @@ def state(
148147
"queue": job.queue,
149148
"state": state,
150149
}
151-
if DualStatsManager is not None:
150+
if AIRFLOW_V_3_2_PLUS:
152151
DualStatsManager.incr("edge_worker.ti.finish", tags={}, legacy_name_tags=tags)
153152
else:
154153
Stats.incr(

providers/edge3/src/airflow/providers/edge3/worker_api/routes/worker.py

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,7 @@
2727
from airflow.api_fastapi.common.router import AirflowRouter
2828
from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
2929
from airflow.providers.common.compat.sdk import Stats, timezone
30-
31-
try:
32-
from airflow.sdk.observability.stats import DualStatsManager
33-
except ImportError:
34-
DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat
30+
from airflow.providers.common.compat.version_compat import AIRFLOW_V_3_2_PLUS
3531
from airflow.providers.edge3.models.edge_worker import EdgeWorkerModel, EdgeWorkerState, set_metrics
3632
from airflow.providers.edge3.worker_api.auth import jwt_token_authorization_rest
3733
from airflow.providers.edge3.worker_api.datamodels import (
@@ -217,7 +213,9 @@ def set_state(
217213
worker.sysinfo = json.dumps(body.sysinfo)
218214
worker.last_update = timezone.utcnow()
219215
session.commit()
220-
if DualStatsManager is not None:
216+
if AIRFLOW_V_3_2_PLUS:
217+
from airflow.sdk.observability.stats import DualStatsManager
218+
221219
DualStatsManager.incr(
222220
"edge_worker.heartbeat_count",
223221
1,

providers/edge3/tests/unit/edge3/executors/test_edge_executor.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,8 +33,9 @@
3333
from airflow.utils.state import TaskInstanceState
3434

3535
from tests_common.test_utils.config import conf_vars
36+
from tests_common.test_utils.version_compat import AIRFLOW_V_3_2_PLUS
3637

37-
try:
38+
if AIRFLOW_V_3_2_PLUS:
3839
from airflow.sdk._shared.observability.metrics.dual_stats_manager import DualStatsManager # noqa: F401
3940

4041
stats_reference = "airflow.sdk._shared.observability.metrics.dual_stats_manager.DualStatsManager"
@@ -48,7 +49,7 @@
4849
},
4950
}
5051
expected_call_count = 1
51-
except ImportError:
52+
else:
5253
from airflow.providers.common.compat.sdk import Stats
5354

5455
stats_reference = f"{Stats.__module__}.Stats"

providers/edge3/tests/unit/edge3/worker_api/routes/test_jobs.py

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,8 @@
2727
from airflow.utils.session import create_session
2828
from airflow.utils.state import TaskInstanceState
2929

30+
from tests_common.test_utils.version_compat import AIRFLOW_V_3_2_PLUS
31+
3032
if TYPE_CHECKING:
3133
from sqlalchemy.orm import Session
3234

@@ -35,7 +37,7 @@
3537
RUN_ID = "manual__2024-11-24T21:03:01+01:00"
3638
QUEUE = "test"
3739

38-
try:
40+
if AIRFLOW_V_3_2_PLUS:
3941
from airflow.sdk._shared.observability.metrics.dual_stats_manager import DualStatsManager # noqa: F401
4042

4143
stats_reference = "airflow.sdk._shared.observability.metrics.dual_stats_manager.DualStatsManager"
@@ -49,7 +51,7 @@
4951
},
5052
}
5153
expected_call_count = 1
52-
except ImportError:
54+
else:
5355
from airflow.providers.common.compat.sdk import Stats
5456

5557
stats_reference = f"{Stats.__module__}.Stats"

providers/openlineage/src/airflow/providers/openlineage/plugins/adapter.py

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -36,11 +36,7 @@
3636
)
3737

3838
from airflow.providers.common.compat.sdk import Stats, conf as airflow_conf
39-
40-
try:
41-
from airflow.sdk.observability.stats import DualStatsManager
42-
except ImportError:
43-
DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat
39+
from airflow.providers.common.compat.version_compat import AIRFLOW_V_3_2_PLUS
4440
from airflow.providers.openlineage import __version__ as OPENLINEAGE_PROVIDER_VERSION, conf
4541
from airflow.providers.openlineage.utils.utils import (
4642
OpenLineageRedactor,
@@ -163,7 +159,9 @@ def emit(self, event: RunEvent):
163159

164160
try:
165161
with ExitStack() as stack:
166-
if DualStatsManager is not None:
162+
if AIRFLOW_V_3_2_PLUS:
163+
from airflow.sdk.observability.stats import DualStatsManager
164+
167165
stack.enter_context(
168166
DualStatsManager.timer(
169167
"ol.emit.attempts",

providers/openlineage/tests/unit/openlineage/plugins/test_adapter.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,11 +60,11 @@
6060
from tests_common.test_utils.taskinstance import create_task_instance
6161
from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS, AIRFLOW_V_3_2_PLUS
6262

63-
try:
63+
if AIRFLOW_V_3_2_PLUS:
6464
from airflow.sdk._shared.observability.metrics.dual_stats_manager import DualStatsManager # noqa: F401
6565

6666
stats_reference = "airflow.sdk._shared.observability.metrics.dual_stats_manager.DualStatsManager"
67-
except ImportError:
67+
else:
6868
stats_reference = "airflow.providers.openlineage.plugins.adapter.Stats"
6969

7070

0 commit comments

Comments
 (0)