This is an automated email from the ASF dual-hosted git repository.
shahar1 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 232c264968f Deprecate Dataproc ClusterGenerator helper class (#71427)
232c264968f is described below
commit 232c264968f3b57918d8c7db85cddbcfbfae3f7b
Author: olegkachur-e <[email protected]>
AuthorDate: Sat Oct 3 07:06:54 2026 +0000
Deprecate Dataproc ClusterGenerator helper class (#71427)
- Becomes obsolete following the deprecation of keyword arguments in
DataprocCreateClusterOperator.
- Direct configuration using 'cluster_config'/'virtual_cluster_config' is
preferred in modern workflows.
- Reduces maintainer overhead by removing legacy helper code.
---
providers/google/docs/operators/cloud/dataproc.rst | 14 +-
.../providers/google/cloud/operators/dataproc.py | 19 ++-
.../dataproc/example_dataproc_batch_persistent.py | 32 +++--
.../dataproc/example_dataproc_cluster_generator.py | 156 ---------------------
.../unit/google/cloud/operators/test_dataproc.py | 7 +-
5 files changed, 40 insertions(+), 188 deletions(-)
diff --git a/providers/google/docs/operators/cloud/dataproc.rst
b/providers/google/docs/operators/cloud/dataproc.rst
index 115b1001eb9..183a6b10443 100644
--- a/providers/google/docs/operators/cloud/dataproc.rst
+++ b/providers/google/docs/operators/cloud/dataproc.rst
@@ -188,16 +188,12 @@ You can use deferrable mode for this action in order to
run the operator asynchr
Generating Cluster Config
^^^^^^^^^^^^^^^^^^^^^^^^^
-You can also generate **CLUSTER_CONFIG** using functional API,
-this could be easily done using **make()** of
-:class:`~airflow.providers.google.cloud.operators.dataproc.ClusterGenerator`
-You can generate and use config as followed:
-.. exampleinclude::
/../../google/tests/system/google/cloud/dataproc/example_dataproc_cluster_generator.py
- :language: python
- :dedent: 0
- :start-after: [START
how_to_cloud_dataproc_create_cluster_generate_cluster_config]
- :end-before: [END
how_to_cloud_dataproc_create_cluster_generate_cluster_config]
+.. warning::
+ **Deprecated:** The
:class:`~airflow.providers.google.cloud.operators.dataproc.ClusterGenerator`
+ class is deprecated and will be removed after September 1, 2027. Please
pass
+ ``cluster_config`` or ``virtual_cluster_config`` dictionaries directly as
arguments to the
+ ``DataprocCreateClusterOperator``.
Diagnose a cluster
------------------
diff --git
a/providers/google/src/airflow/providers/google/cloud/operators/dataproc.py
b/providers/google/src/airflow/providers/google/cloud/operators/dataproc.py
index e41e9b8550c..6ef81c4bf15 100644
--- a/providers/google/src/airflow/providers/google/cloud/operators/dataproc.py
+++ b/providers/google/src/airflow/providers/google/cloud/operators/dataproc.py
@@ -22,7 +22,6 @@ from __future__ import annotations
import inspect
import re
import time
-import warnings
from collections.abc import MutableSequence, Sequence
from dataclasses import dataclass
from datetime import datetime, timedelta
@@ -62,6 +61,7 @@ from airflow.providers.google.cloud.triggers.dataproc import (
DataprocSubmitTrigger,
)
from airflow.providers.google.cloud.utils.dataproc import DataprocOperationType
+from airflow.providers.google.common.deprecated import deprecated
from airflow.providers.google.common.hooks.base_google import
PROVIDE_PROJECT_ID
from airflow.triggers.base import StartTriggerArgs
@@ -117,6 +117,14 @@ class InstanceFlexibilityPolicy:
instance_selection_list: list[InstanceSelection]
+@deprecated(
+ planned_removal_date="September 1, 2027",
+ reason="Since passing cluster parameters by keyword in the
DataprocCreateClusterOperator is scheduled "
+ "for deletion, there is no longer a need to maintain the
'ClusterGenerator' helper class.",
+ instructions="Please pass 'cluster_config' or 'virtual_cluster_config'
directly as "
+ "DataprocCreateClusterOperator arguments.",
+ category=AirflowProviderDeprecationWarning,
+)
class ClusterGenerator:
"""
Create a new Dataproc Cluster.
@@ -712,16 +720,7 @@ class
DataprocCreateClusterOperator(GoogleCloudBaseOperator):
polling_interval_seconds: int = 10,
**kwargs,
) -> None:
- # TODO: remove one day
if cluster_config is None and virtual_cluster_config is None:
- warnings.warn(
- f"Passing cluster parameters by keywords to
`{type(self).__name__}` will be deprecated. "
- "Please provide cluster_config object using `cluster_config`
parameter. "
- "You can use
`airflow.dataproc.ClusterGenerator.generate_cluster` "
- "method to obtain cluster object. Planned removal date:
October 5, 2026.",
- AirflowProviderDeprecationWarning,
- stacklevel=2,
- )
# Remove result of apply defaults
if "params" in kwargs:
del kwargs["params"]
diff --git
a/providers/google/tests/system/google/cloud/dataproc/example_dataproc_batch_persistent.py
b/providers/google/tests/system/google/cloud/dataproc/example_dataproc_batch_persistent.py
index c13f4e0062c..521753bcdca 100644
---
a/providers/google/tests/system/google/cloud/dataproc/example_dataproc_batch_persistent.py
+++
b/providers/google/tests/system/google/cloud/dataproc/example_dataproc_batch_persistent.py
@@ -27,7 +27,6 @@ from google.api_core.retry import Retry
from airflow.models.dag import DAG
from airflow.providers.google.cloud.operators.dataproc import (
- ClusterGenerator,
DataprocCreateBatchOperator,
DataprocCreateClusterOperator,
DataprocDeleteBatchOperator,
@@ -53,17 +52,26 @@ CLUSTER_NAME_FULL = CLUSTER_NAME_BASE +
f"-{ENV_ID}".replace("_", "-")
CLUSTER_NAME = CLUSTER_NAME_BASE if len(CLUSTER_NAME_FULL) >= 33 else
CLUSTER_NAME_FULL
BATCH_ID = f"batch-{ENV_ID}-{DAG_ID}".replace("_", "-")
-CLUSTER_GENERATOR_CONFIG_FOR_PHS = ClusterGenerator(
- project_id=PROJECT_ID,
- region=REGION,
- master_machine_type="n1-standard-4",
- worker_machine_type="n1-standard-4",
- num_workers=0,
- properties={
- "spark:spark.history.fs.logDirectory": f"gs://{BUCKET_NAME}",
+CLUSTER_CONFIG_FOR_PHS = {
+ "master_config": {
+ "num_instances": 1,
+ "machine_type_uri": "n1-standard-4",
},
- enable_component_gateway=True,
-).make()
+ "worker_config": {
+ "num_instances": 0,
+ "machine_type_uri": "n1-standard-4",
+ },
+ "software_config": {
+ "properties": {
+ "dataproc:dataproc.allow.zero.workers": "true",
+ "spark:spark.history.fs.logDirectory": f"gs://{BUCKET_NAME}",
+ },
+ },
+ "endpoint_config": {
+ "enable_http_port_access": True,
+ },
+}
+
BATCH_CONFIG_WITH_PHS = {
"spark_batch": {
"jar_file_uris":
["file:///usr/lib/spark/examples/jars/spark-examples.jar"],
@@ -94,7 +102,7 @@ with DAG(
create_cluster = DataprocCreateClusterOperator(
task_id="create_cluster_for_phs",
project_id=PROJECT_ID,
- cluster_config=CLUSTER_GENERATOR_CONFIG_FOR_PHS,
+ cluster_config=CLUSTER_CONFIG_FOR_PHS,
region=REGION,
cluster_name=CLUSTER_NAME,
retry=Retry(maximum=100.0, initial=10.0, multiplier=1.0),
diff --git
a/providers/google/tests/system/google/cloud/dataproc/example_dataproc_cluster_generator.py
b/providers/google/tests/system/google/cloud/dataproc/example_dataproc_cluster_generator.py
deleted file mode 100644
index d7c639a1af0..00000000000
---
a/providers/google/tests/system/google/cloud/dataproc/example_dataproc_cluster_generator.py
+++ /dev/null
@@ -1,156 +0,0 @@
-#
-# 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.
-"""
-Example Airflow DAG testing Managed Spark
-operators for managing a cluster and submitting jobs.
-"""
-
-from __future__ import annotations
-
-import os
-from datetime import datetime
-
-from google.api_core.retry import Retry
-
-from airflow.models.dag import DAG
-from airflow.providers.google.cloud.operators.dataproc import (
- ClusterGenerator,
- DataprocCreateClusterOperator,
- DataprocDeleteClusterOperator,
-)
-from airflow.providers.google.cloud.operators.gcs import (
- GCSCreateBucketOperator,
- GCSDeleteBucketOperator,
- GCSSynchronizeBucketsOperator,
-)
-
-try:
- from airflow.sdk import TriggerRule
-except ImportError:
- # Compatibility for Airflow < 3.1
- from airflow.utils.trigger_rule import TriggerRule # type:
ignore[no-redef,attr-defined]
-
-from system.google import DEFAULT_GCP_SYSTEM_TEST_PROJECT_ID
-
-ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
-DAG_ID = "dataproc_cluster_generation"
-PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT") or
DEFAULT_GCP_SYSTEM_TEST_PROJECT_ID
-
-BUCKET_NAME = f"bucket_{DAG_ID}_{ENV_ID}"
-RESOURCE_DATA_BUCKET = "airflow-system-tests-resources"
-INIT_FILE = "pip-install.sh"
-GCS_INIT_FILE = f"gs://{RESOURCE_DATA_BUCKET}/dataproc/{INIT_FILE}"
-
-CLUSTER_NAME_BASE = f"cluster-{DAG_ID}".replace("_", "-")
-CLUSTER_NAME_FULL = CLUSTER_NAME_BASE + f"-{ENV_ID}".replace("_", "-")
-CLUSTER_NAME = CLUSTER_NAME_BASE if len(CLUSTER_NAME_FULL) >= 33 else
CLUSTER_NAME_FULL
-
-REGION = "us-east4"
-ZONE = "us-east4-a"
-
-# Cluster definition: Generating Cluster Config for
DataprocCreateClusterOperator
-# [START how_to_cloud_dataproc_create_cluster_generate_cluster_config]
-CLUSTER_GENERATOR_CONFIG = ClusterGenerator(
- project_id=PROJECT_ID,
- zone=ZONE,
- master_machine_type="n1-standard-4",
- master_disk_size=32,
- worker_machine_type="n1-standard-4",
- worker_disk_size=32,
- num_workers=2,
- storage_bucket=BUCKET_NAME,
- init_actions_uris=[GCS_INIT_FILE],
- metadata={"PIP_PACKAGES": "pyyaml requests pandas openpyxl"},
- num_preemptible_workers=1,
- preemptibility="PREEMPTIBLE",
- internal_ip_only=False,
- cluster_tier="CLUSTER_TIER_STANDARD",
- cluster_type="STANDARD",
- engine="DEFAULT",
-).make()
-
-# [END how_to_cloud_dataproc_create_cluster_generate_cluster_config]
-
-
-with DAG(
- DAG_ID,
- schedule="@once",
- start_date=datetime(2021, 1, 1),
- catchup=False,
- tags=["example", "managed-spark"],
-) as dag:
- create_bucket = GCSCreateBucketOperator(
- task_id="create_bucket", bucket_name=BUCKET_NAME, project_id=PROJECT_ID
- )
-
- move_init_file = GCSSynchronizeBucketsOperator(
- task_id="move_init_file",
- source_bucket=RESOURCE_DATA_BUCKET,
- source_object="dataproc",
- destination_bucket=BUCKET_NAME,
- destination_object="dataproc",
- recursive=True,
- )
-
- # [START
how_to_cloud_dataproc_create_cluster_generate_cluster_config_operator]
-
- create_dataproc_cluster = DataprocCreateClusterOperator(
- task_id="create_dataproc_cluster",
- cluster_name=CLUSTER_NAME,
- project_id=PROJECT_ID,
- region=REGION,
- cluster_config=CLUSTER_GENERATOR_CONFIG,
- retry=Retry(maximum=100.0, initial=10.0, multiplier=1.0),
- num_retries_if_resource_is_not_ready=3,
- )
-
- # [END
how_to_cloud_dataproc_create_cluster_generate_cluster_config_operator]
-
- delete_cluster = DataprocDeleteClusterOperator(
- task_id="delete_cluster",
- project_id=PROJECT_ID,
- cluster_name=CLUSTER_NAME,
- region=REGION,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- delete_bucket = GCSDeleteBucketOperator(
- task_id="delete_bucket", bucket_name=BUCKET_NAME,
trigger_rule=TriggerRule.ALL_DONE
- )
-
- (
- # TEST SETUP
- create_bucket
- >> move_init_file
- # TEST BODY
- >> create_dataproc_cluster
- # TEST TEARDOWN
- >> [delete_cluster, delete_bucket]
- )
-
- from tests_common.test_utils.watcher import watcher
-
- # This test needs watcher in order to properly mark success/failure
- # when "teardown" task with trigger rule is part of the DAG
- list(dag.tasks) >> watcher()
-
-
-from tests_common.test_utils.system_tests import get_test_run # noqa: E402
-
-# Needed to run the example DAG with pytest (see:
contributing-docs/testing/system_tests.rst)
-test_run = get_test_run(dag)
diff --git
a/providers/google/tests/unit/google/cloud/operators/test_dataproc.py
b/providers/google/tests/unit/google/cloud/operators/test_dataproc.py
index 7f0f13f77bd..c1cb55f566f 100644
--- a/providers/google/tests/unit/google/cloud/operators/test_dataproc.py
+++ b/providers/google/tests/unit/google/cloud/operators/test_dataproc.py
@@ -632,6 +632,7 @@ class DataprocClusterTestBase(DataprocTestBase):
]
[email protected]("ignore::airflow.exceptions.AirflowProviderDeprecationWarning")
class TestsClusterGenerator:
def test_image_version(self):
with pytest.raises(ValueError, match="custom_image and image_version"):
@@ -982,6 +983,10 @@ class TestsClusterGenerator:
cluster = generator.make()
assert cluster["engine"] == "DEFAULT"
+ def test_deprecation_warning(self):
+ with pytest.warns(AirflowProviderDeprecationWarning,
match="ClusterGenerator"):
+ ClusterGenerator(project_id=GCP_PROJECT)
+
class TestDataprocCreateClusterOperator(DataprocClusterTestBase):
def test_deprecation_warning(self):
@@ -994,7 +999,7 @@ class
TestDataprocCreateClusterOperator(DataprocClusterTestBase):
num_workers=2,
zone="zone",
)
- assert_warning("Passing cluster parameters by keywords", warnings)
+ assert_warning("Since passing cluster parameters by keyword", warnings)
assert op.project_id == GCP_PROJECT
assert op.cluster_name == "cluster_name"