Skip to content

Commit 01bc19c

Browse files
Oleg Kachurolegkachur-e
authored andcommitted
Revert "Suspend Apache Beam Provider due to grpcio limitation (#61926)"
This reverts commit 917abea. - Remove hacks regarding the beam provider suspension, as they are not relevant anymore. ISSUE: #66551
1 parent 80f1ab4 commit 01bc19c

13 files changed

Lines changed: 20 additions & 86 deletions

File tree

airflow-core/tests/unit/always/test_example_dags.py

Lines changed: 0 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -69,17 +69,6 @@
6969
# Ray uses pydantic v1 internally, which fails to infer types in Python 3.14.
7070
# TODO: remove once ray releases a version with Python 3.14 support.
7171
"providers/google/tests/system/google/cloud/ray/example_ray_job.py",
72-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_go.py",
73-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_java_streaming.py",
74-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_native_java.py",
75-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_native_python.py",
76-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_native_python_async.py",
77-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_pipeline.py",
78-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_pipeline_streaming.py",
79-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_sensors_deferrable.py",
80-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_streaming_python.py",
81-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_template.py",
82-
"providers/google/tests/system/google/cloud/dataflow/example_dataflow_yaml.py",
8372
"providers/google/tests/system/google/cloud/gcs/example_firestore.py",
8473
)
8574

dev/breeze/tests/test_packages.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -118,12 +118,12 @@ def test_get_removed_providers():
118118

119119
def test_get_suspended_provider_ids():
120120
# Modify it every time we suspend/resume provider
121-
assert get_suspended_provider_ids() == ["apache.beam"]
121+
assert get_suspended_provider_ids() == []
122122

123123

124124
def test_get_suspended_provider_folders():
125125
# Modify it every time we suspend/resume provider
126-
assert get_suspended_provider_folders() == ["apache/beam"]
126+
assert get_suspended_provider_folders() == []
127127

128128

129129
@pytest.mark.parametrize(

devel-common/src/tests_common/test_utils/providers.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,7 +68,7 @@ def get_provider_min_airflow_version(provider_name: str) -> tuple[int, ...]:
6868

6969

7070
# Ignore module import errors for any suspended provider paths used in example dags
71-
IGNORE_MODULE_IMPORT_ERRORS: list[str] = ["airflow.providers.apache.beam"]
71+
IGNORE_MODULE_IMPORT_ERRORS: list[str] = []
7272

7373

7474
def get_suspended_providers_folders() -> list[str]:

providers/apache/beam/provider.yaml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,9 +21,9 @@ name: Apache Beam
2121
description: |
2222
`Apache Beam <https://beam.apache.org/>`__.
2323
24-
state: suspended
25-
source-date-epoch: 1768334134
24+
state: ready
2625
lifecycle: production
26+
source-date-epoch: 1772063928
2727
# Note that those versions are maintained by release manager - do not update them manually
2828
# with the exception of case where other provider in sources has >= new provider version.
2929
# In such case adding >= NEW_VERSION and bumping to NEW_VERSION in a provider have

providers/google/pyproject.toml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,9 @@ dependencies = [
149149
# The optional dependencies should be modified in place in the generated file
150150
# Any change in the dependencies is preserved when the file is regenerated
151151
[project.optional-dependencies]
152+
"apache.beam" = [
153+
"apache-airflow-providers-apache-beam>=6.2.2",
154+
]
152155
"cncf.kubernetes" = [
153156
"apache-airflow-providers-cncf-kubernetes>=10.1.0",
154157
]
@@ -216,6 +219,7 @@ dev = [
216219
"apache-airflow-task-sdk",
217220
"apache-airflow-devel-common",
218221
"apache-airflow-providers-amazon",
222+
"apache-airflow-providers-apache-beam",
219223
"apache-airflow-providers-apache-cassandra",
220224
"apache-airflow-providers-cncf-kubernetes",
221225
"apache-airflow-providers-common-compat",

providers/google/tests/unit/google/cloud/hooks/test_dataflow.py

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,6 @@
2828
from uuid import UUID
2929

3030
import pytest
31-
32-
# TODO: Remove below skip once beam provider changed to ready state
33-
pytest.importorskip("apache-beam", reason="apache-beam package suspended due to grpcio limitation")
34-
35-
3631
from google.cloud.dataflow_v1beta3 import (
3732
GetJobMetricsRequest,
3833
GetJobRequest,

providers/google/tests/unit/google/cloud/operators/test_dataflow.py

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -21,10 +21,6 @@
2121

2222
import httplib2
2323
import pytest
24-
25-
# TODO: Remove below skip once beam provider changed to ready state
26-
pytest.importorskip("apache-beam", reason="apache-beam package suspended due to grpcio limitation")
27-
2824
from googleapiclient.errors import HttpError
2925

3026
from airflow.providers.common.compat.sdk import AirflowException

providers/google/tests/unit/google/cloud/operators/test_datapipeline.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,9 +21,6 @@
2121

2222
import pytest
2323

24-
# TODO: Remove below skip once beam provider changed to ready state
25-
pytest.importorskip("apache-beam", reason="apache-beam package suspended due to grpcio limitation")
26-
2724
from airflow.providers.google.cloud.operators.dataflow import (
2825
DataflowCreatePipelineOperator,
2926
DataflowRunPipelineOperator,

providers/google/tests/unit/google/cloud/sensors/test_dataflow.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,9 +21,6 @@
2121

2222
import pytest
2323

24-
# TODO: Remove below skip once beam provider changed to ready state
25-
pytest.importorskip("apache-beam", reason="apache-beam package suspended due to grpcio limitation")
26-
2724
from airflow.providers.common.compat.sdk import AirflowException, TaskDeferred
2825
from airflow.providers.google.cloud.hooks.dataflow import DataflowJobStatus
2926
from airflow.providers.google.cloud.sensors.dataflow import (

providers/google/tests/unit/google/cloud/triggers/test_dataflow.py

Lines changed: 3 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -19,58 +19,14 @@
1919

2020
import asyncio
2121
import logging
22-
import sys
23-
import types
2422
from unittest import mock
2523

2624
import pytest
2725
from google.api_core.exceptions import ServiceUnavailable
2826
from google.cloud.dataflow_v1beta3 import Job, JobState, JobType
2927

30-
31-
# While the apache-beam provider is suspended (#61926),
32-
# `airflow.providers.google.cloud.hooks.dataflow` cannot be imported because it
33-
# does `from airflow.providers.apache.beam.hooks.beam import ...` at module
34-
# scope. Probe-import the real module first (which also triggers Airflow's
35-
# provider-entry-point discovery while the real `airflow.providers.apache`
36-
# namespace package is still intact) and only fall back to stubs when the
37-
# probe fails. The stubs deliberately cover only the missing beam-specific
38-
# subtree — we never overwrite `airflow.providers.apache` itself, otherwise
39-
# sibling providers such as `airflow.providers.apache.pig` would be unreachable
40-
# via the namespace package's `__path__`.
41-
def _beam_module_stubs() -> dict[str, types.ModuleType]:
42-
beam_package = types.ModuleType("airflow.providers.apache.beam")
43-
beam_package.__path__ = []
44-
hooks_package = types.ModuleType("airflow.providers.apache.beam.hooks")
45-
hooks_package.__path__ = []
46-
47-
class BeamRunnerType:
48-
DataflowRunner = "DataflowRunner"
49-
50-
class BeamModule(types.ModuleType):
51-
BeamHook: object
52-
BeamRunnerType: object
53-
beam_options_to_args: object
54-
55-
beam_module = BeamModule("airflow.providers.apache.beam.hooks.beam")
56-
beam_module.BeamHook = mock.MagicMock()
57-
beam_module.BeamRunnerType = BeamRunnerType
58-
beam_module.beam_options_to_args = mock.MagicMock(return_value=[])
59-
60-
return {
61-
"airflow.providers.apache.beam": beam_package,
62-
"airflow.providers.apache.beam.hooks": hooks_package,
63-
"airflow.providers.apache.beam.hooks.beam": beam_module,
64-
}
65-
66-
67-
try:
68-
import airflow.providers.apache.beam.hooks.beam # noqa: F401
69-
except ImportError:
70-
sys.modules.update(_beam_module_stubs())
71-
72-
from airflow.providers.google.cloud.hooks.dataflow import DataflowJobStatus # noqa: E402
73-
from airflow.providers.google.cloud.triggers.dataflow import ( # noqa: E402
28+
from airflow.providers.google.cloud.hooks.dataflow import DataflowJobStatus
29+
from airflow.providers.google.cloud.triggers.dataflow import (
7430
DataflowJobAutoScalingEventTrigger,
7531
DataflowJobMessagesTrigger,
7632
DataflowJobMetricsTrigger,
@@ -79,7 +35,7 @@ class BeamModule(types.ModuleType):
7935
DataflowStartYamlJobTrigger,
8036
TemplateJobStartTrigger,
8137
)
82-
from airflow.triggers.base import TriggerEvent # noqa: E402
38+
from airflow.triggers.base import TriggerEvent
8339

8440
PROJECT_ID = "test-project-id"
8541
JOB_ID = "test_job_id_2012-12-23-10:00"

0 commit comments

Comments
 (0)