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"

Reply via email to