diff --git a/docs/apache-airflow-providers-google/operators/cloud/automl.rst b/docs/apache-airflow-providers-google/operators/cloud/automl.rst
index b93b8d8cf75ff..4eb461409fa3c 100644
--- a/docs/apache-airflow-providers-google/operators/cloud/automl.rst
+++ b/docs/apache-airflow-providers-google/operators/cloud/automl.rst
@@ -16,7 +16,14 @@
under the License.
Google Cloud AutoML Operators
-=======================================
+=============================
+
+.. warning::
+ The AutoML API is deprecated. Planned removal date is September 30, 2025, but some operators might be deleted
+ earlier, according to the docs and deprecation warnings!
+ The replacement suggestions can be found in the deprecation warnings or in the doc below.
+ Please note that AutoML for translation API functionality has been moved to the Advanced Translation service,
+ the operators can be found at ``airflow.providers.google.cloud.operators.translate`` module.
The `Google Cloud AutoML `__
makes the power of machine learning available to you even if you have limited knowledge
@@ -41,27 +48,15 @@ To create a Google AutoML dataset you can use
:class:`~airflow.providers.google.cloud.operators.automl.AutoMLCreateDatasetOperator`.
The operator returns dataset id in :ref:`XCom ` under ``dataset_id`` key.
-.. warning::
- This operator is deprecated when running for text, video and vision prediction and will be removed soon.
- All the functionality of legacy AutoML Natural Language, Vision, Video Intelligence and new features are
- available on the Vertex AI platform. Please use
- :class:`~airflow.providers.google.cloud.operators.vertex_ai.dataset.CreateDatasetOperator`
-
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_dataset.py
- :language: python
- :dedent: 4
- :start-after: [START howto_operator_automl_create_dataset]
- :end-before: [END howto_operator_automl_create_dataset]
+This operator is deprecated when running for text, video and vision prediction and will be removed after September 30, 2025.
+All the functionality of legacy AutoML Natural Language, Vision, Video Intelligence and new features are
+available on the Vertex AI platform. Please use
+:class:`~airflow.providers.google.cloud.operators.vertex_ai.dataset.CreateDatasetOperator`
+:class:`~airflow.providers.google.cloud.operators.translate.TranslateCreateDatasetOperator`.
After creating a dataset you can use it to import some data using
:class:`~airflow.providers.google.cloud.operators.automl.AutoMLImportDataOperator`.
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_dataset.py
- :language: python
- :dedent: 4
- :start-after: [START howto_operator_automl_import_data]
- :end-before: [END howto_operator_automl_import_data]
-
To update dataset you can use
:class:`~airflow.providers.google.cloud.operators.automl.AutoMLTablesUpdateDatasetOperator`.
@@ -105,28 +100,14 @@ The operator will wait for the operation to complete. Additionally the operator
returns the id of model in :ref:`XCom ` under ``model_id`` key.
.. warning::
- This operator is deprecated when running for text, video and vision prediction and will be removed soon.
+ This operator is deprecated when running for text, video and vision prediction and will be removed after September 30, 2025.
All the functionality of legacy AutoML Natural Language, Vision, Video Intelligence and new features are
available on the Vertex AI platform. Please use
+ :class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLTabularTrainingJobOperator`,
+ :class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLVideoTrainingJobOperator`,
+ :class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLImageTrainingJobOperator`,
:class:`~airflow.providers.google.cloud.operators.vertex_ai.generative_model.SupervisedFineTuningTrainOperator`,
- :class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLImageTrainingJobOperator` or
- :class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLVideoTrainingJobOperator`.
-
- You can find example on how to use VertexAI operators for AutoML Vision classification here:
-
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_vision_classification.py
- :language: python
- :dedent: 4
- :start-after: [START howto_cloud_create_image_classification_training_job_operator]
- :end-before: [END howto_cloud_create_image_classification_training_job_operator]
-
-Example on how to use VertexAI operators for AutoML Video Intelligence classification you can find here:
-
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_video_classification.py
- :language: python
- :dedent: 4
- :start-after: [START howto_cloud_create_video_classification_training_job_operator]
- :end-before: [END howto_cloud_create_video_classification_training_job_operator]
+ :class:`~airflow.providers.google.cloud.operators.translate.TranslateCreateModelOperator`.
When running Vertex AI Operator for training data, please ensure that your data is correctly stored in Vertex AI
datasets. To create and import data to the dataset please use
@@ -134,18 +115,17 @@ datasets. To create and import data to the dataset please use
and
:class:`~airflow.providers.google.cloud.operators.vertex_ai.dataset.ImportDataOperator`
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_translation.py
- :language: python
- :dedent: 4
- :start-after: [START howto_operator_automl_create_model]
- :end-before: [END howto_operator_automl_create_model]
+For the AutoML translation please use the
+:class:`~airflow.providers.google.cloud.operators.translate.TranslateTextOperator`
+or
+:class:`~airflow.providers.google.cloud.operators.translate.TranslateTextBatchOperator`.
To get existing model one can use
:class:`~airflow.providers.google.cloud.operators.automl.AutoMLGetModelOperator`.
This operator deprecated for tables, video intelligence, vision and natural language is deprecated
and will be removed after 31.03.2024. Please use
-:class:`airflow.providers.google.cloud.operators.vertex_ai.model_service.GetModelOperator` instead.
+:class:`~airflow.providers.google.cloud.operators.vertex_ai.model_service.GetModelOperator` instead.
You can find example on how to use VertexAI operators here:
.. exampleinclude:: /../../providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_model_service.py
@@ -193,14 +173,10 @@ To obtain predictions from Google Cloud AutoML model you can use
:class:`~airflow.providers.google.cloud.operators.automl.AutoMLBatchPredictOperator`. In the first case
the model must be deployed.
-.. exampleinclude:: /../../providers/tests/system/google/cloud/automl/example_automl_translation.py
- :language: python
- :dedent: 4
- :start-after: [START howto_operator_prediction]
- :end-before: [END howto_operator_prediction]
Th :class:`~airflow.providers.google.cloud.operators.automl.AutoMLBatchPredictOperator` deprecated for tables,
-video intelligence, vision and natural language is deprecated and will be removed after 31.03.2024. Please use
+video intelligence, vision and natural language is deprecated and will be removed after 31.03.2024.
+Please use
:class:`airflow.providers.google.cloud.operators.vertex_ai.batch_prediction_job.CreateBatchPredictionJobOperator`,
:class:`airflow.providers.google.cloud.operators.vertex_ai.batch_prediction_job.GetBatchPredictionJobOperator`,
:class:`airflow.providers.google.cloud.operators.vertex_ai.batch_prediction_job.ListBatchPredictionJobsOperator`,
@@ -238,7 +214,9 @@ of datasets ids in :ref:`XCom ` under ``dataset_id_list`` key.
This operator deprecated for tables, video intelligence, vision and natural language is deprecated
and will be removed after 31.03.2024. Please use
-:class:`airflow.providers.google.cloud.operators.vertex_ai.dataset.ListDatasetsOperator` instead.
+:class:`~airflow.providers.google.cloud.operators.vertex_ai.dataset.ListDatasetsOperator`,
+:class:`~airflow.providers.google.cloud.operators.translate.TranslateDatasetsListOperator`
+instead.
You can find example on how to use VertexAI operators here:
.. exampleinclude:: /../../providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_dataset.py
diff --git a/docs/apache-airflow-providers-google/operators/cloud/vertex_ai.rst b/docs/apache-airflow-providers-google/operators/cloud/vertex_ai.rst
index 12b86d25d9196..e538241959a69 100644
--- a/docs/apache-airflow-providers-google/operators/cloud/vertex_ai.rst
+++ b/docs/apache-airflow-providers-google/operators/cloud/vertex_ai.rst
@@ -245,6 +245,14 @@ put dataset id to ``dataset_id`` parameter in operator.
:start-after: [START how_to_cloud_vertex_ai_create_auto_ml_image_training_job_operator]
:end-before: [END how_to_cloud_vertex_ai_create_auto_ml_image_training_job_operator]
+To run AutoML image detection training job:
+
+.. exampleinclude:: /../../providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_image_object_detection.py
+ :language: python
+ :dedent: 4
+ :start-after: [START how_to_cloud_vertex_ai_create_auto_ml_image_object_detection_training_job_operator]
+ :end-before: [END how_to_cloud_vertex_ai_create_auto_ml_image_object_detection_training_job_operator]
+
How to run AutoML Tabular Training Job
:class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLTabularTrainingJobOperator`
@@ -279,6 +287,15 @@ version of existing Model instead of new Model created in Model Registry. This c
:start-after: [START how_to_cloud_vertex_ai_create_auto_ml_video_training_job_v2_operator]
:end-before: [END how_to_cloud_vertex_ai_create_auto_ml_video_training_job_v2_operator]
+Also you can use vertex_ai AutoML model for video tracking.
+
+.. exampleinclude:: /../../providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_video_tracking.py
+ :language: python
+ :dedent: 4
+ :start-after: [START how_to_cloud_vertex_ai_create_auto_ml_video_tracking_job_operator]
+ :end-before: [END how_to_cloud_vertex_ai_create_auto_ml_video_tracking_job_operator]
+
+
You can get a list of AutoML Training Jobs using
:class:`~airflow.providers.google.cloud.operators.vertex_ai.auto_ml.ListAutoMLTrainingJobOperator`.
diff --git a/providers/src/airflow/providers/google/cloud/hooks/automl.py b/providers/src/airflow/providers/google/cloud/hooks/automl.py
index 83c61e8ddfa37..ba13b1e4edb77 100644
--- a/providers/src/airflow/providers/google/cloud/hooks/automl.py
+++ b/providers/src/airflow/providers/google/cloud/hooks/automl.py
@@ -43,8 +43,9 @@
PredictResponse,
)
-from airflow.exceptions import AirflowException
+from airflow.exceptions import AirflowException, AirflowProviderDeprecationWarning
from airflow.providers.google.common.consts import CLIENT_INFO
+from airflow.providers.google.common.deprecated import deprecated
from airflow.providers.google.common.hooks.base_google import PROVIDE_PROJECT_ID, GoogleBaseHook
if TYPE_CHECKING:
@@ -58,6 +59,12 @@
from google.protobuf.field_mask_pb2 import FieldMask
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.hooks.vertex_ai.auto_ml.AutoMLHook, "
+ "airflow.providers.google.cloud.hooks.translate.TranslateHook",
+ category=AirflowProviderDeprecationWarning,
+)
class CloudAutoMLHook(GoogleBaseHook):
"""
Google Cloud AutoML hook.
diff --git a/providers/src/airflow/providers/google/cloud/operators/automl.py b/providers/src/airflow/providers/google/cloud/operators/automl.py
index b49a1b6b67835..7ef0716615126 100644
--- a/providers/src/airflow/providers/google/cloud/operators/automl.py
+++ b/providers/src/airflow/providers/google/cloud/operators/automl.py
@@ -20,7 +20,6 @@
from __future__ import annotations
import ast
-import warnings
from collections.abc import Sequence
from functools import cached_property
from typing import TYPE_CHECKING, cast
@@ -57,26 +56,15 @@
MetaData = Sequence[tuple[str, str]]
-def _raise_exception_for_deprecated_operator(
- deprecated_class_name: str, alternative_class_names: str | list[str]
-):
- if isinstance(alternative_class_names, str):
- alternative_class_name_str = alternative_class_names
- elif len(alternative_class_names) == 1:
- alternative_class_name_str = alternative_class_names[0]
- else:
- alternative_class_name_str = ", ".join(f"`{cls_name}`" for cls_name in alternative_class_names[:-1])
- alternative_class_name_str += f" or `{alternative_class_names[-1]}`"
-
- raise AirflowException(
- f"{deprecated_class_name} for text, image, and video prediction has been "
- f"deprecated and no longer available. All the functionality of "
- f"legacy AutoML Natural Language, Vision, Video Intelligence and Tables "
- f"and new features are available on the Vertex AI platform. "
- f"Please use {alternative_class_name_str} from Vertex AI."
- )
-
-
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLTabularTrainingJobOperator, "
+ "airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLVideoTrainingJobOperator, "
+ "airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLImageTrainingJobOperator, "
+ "airflow.providers.google.cloud.operators.vertex_ai.generative_model.SupervisedFineTuningTrainOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateCreateModelOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLTrainModelOperator(GoogleCloudBaseOperator):
"""
Creates Google Cloud AutoML model.
@@ -88,6 +76,7 @@ class AutoMLTrainModelOperator(GoogleCloudBaseOperator):
:class:`airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLVideoTrainingJobOperator`,
:class:`airflow.providers.google.cloud.operators.vertex_ai.auto_ml.CreateAutoMLImageTrainingJobOperator`,
:class:`airflow.providers.google.cloud.operators.vertex_ai.generative_model.SupervisedFineTuningTrainOperator`,
+ :class:`airflow.providers.google.cloud.operators.translate.TranslateCreateModelOperator`.
instead.
.. seealso::
@@ -149,17 +138,6 @@ def __init__(
self.impersonation_chain = impersonation_chain
def execute(self, context: Context):
- # Raise exception if running not AutoML Translation prediction job
- if "translation_model_metadata" not in self.model:
- _raise_exception_for_deprecated_operator(
- self.__class__.__name__,
- [
- "CreateAutoMLTabularTrainingJobOperator",
- "CreateAutoMLVideoTrainingJobOperator",
- "CreateAutoMLImageTrainingJobOperator",
- "SupervisedFineTuningTrainOperator",
- ],
- )
hook = CloudAutoMLHook(
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
@@ -195,6 +173,11 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.translate.TranslateTextOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLPredictOperator(GoogleCloudBaseOperator):
"""
Runs prediction operation on Google Cloud AutoML.
@@ -297,20 +280,10 @@ def model(self) -> Model | None:
)
return None
- def _check_model_type(self):
- if not hasattr(self.model, "translation_model_metadata"):
- raise AirflowException(
- "AutoMLPredictOperator for text, image, and video prediction has been deprecated. "
- "Please use endpoint_id param instead of model_id param."
- )
-
def execute(self, context: Context):
if self.model_id is None and self.endpoint_id is None:
raise AirflowException("You must specify model_id or endpoint_id!")
- if self.model_id:
- self._check_model_type()
-
hook = self.hook
if self.model_id:
result = hook.predict(
@@ -461,16 +434,6 @@ def model(self) -> Model:
)
def execute(self, context: Context):
- if not hasattr(self.model, "translation_model_metadata"):
- _raise_exception_for_deprecated_operator(
- self.__class__.__name__,
- [
- "CreateBatchPredictionJobOperator",
- "GetBatchPredictionJobOperator",
- "ListBatchPredictionJobsOperator",
- "DeleteBatchPredictionJobOperator",
- ],
- )
self.log.info("Fetch batch prediction.")
operation = self.hook.batch_predict(
model_id=self.model_id,
@@ -498,13 +461,20 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.dataset.CreateDatasetOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateCreateDatasetOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLCreateDatasetOperator(GoogleCloudBaseOperator):
"""
Creates a Google Cloud AutoML dataset.
AutoMLCreateDatasetOperator for tables, video intelligence, vision and natural language has been
deprecated and no longer available. Please use
- :class:`airflow.providers.google.cloud.operators.vertex_ai.dataset.CreateDatasetOperator` instead.
+ :class:`airflow.providers.google.cloud.operators.vertex_ai.dataset.CreateDatasetOperator`,
+ :class:`airflow.providers.google.cloud.operators.translate.TranslateCreateDatasetOperator` instead.
.. seealso::
For more information on how to use this operator, take a look at the guide:
@@ -565,8 +535,6 @@ def __init__(
self.impersonation_chain = impersonation_chain
def execute(self, context: Context):
- if "translation_dataset_metadata" not in self.dataset:
- _raise_exception_for_deprecated_operator(self.__class__.__name__, "CreateDatasetOperator")
hook = CloudAutoMLHook(
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
@@ -596,6 +564,12 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.dataset.ImportDataOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateImportDataOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLImportDataOperator(GoogleCloudBaseOperator):
"""
Imports data to a Google Cloud AutoML dataset.
@@ -672,7 +646,7 @@ def execute(self, context: Context):
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
)
- dataset: Dataset = hook.get_dataset(
+ hook.get_dataset(
dataset_id=self.dataset_id,
location=self.location,
project_id=self.project_id,
@@ -680,8 +654,6 @@ def execute(self, context: Context):
timeout=self.timeout,
metadata=self.metadata,
)
- if not hasattr(dataset, "translation_dataset_metadata"):
- _raise_exception_for_deprecated_operator(self.__class__.__name__, "ImportDataOperator")
self.log.info("Importing data to dataset...")
operation = hook.import_data(
dataset_id=self.dataset_id,
@@ -704,6 +676,11 @@ def execute(self, context: Context):
)
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ category=AirflowProviderDeprecationWarning,
+ reason="Shutdown of legacy version of AutoML Tables on March 31, 2024.",
+)
class AutoMLTablesListColumnSpecsOperator(GoogleCloudBaseOperator):
"""
Lists column specs in a table.
@@ -787,11 +764,6 @@ def __init__(
self.retry = retry
self.gcp_conn_id = gcp_conn_id
self.impersonation_chain = impersonation_chain
- raise AirflowException(
- "Operator AutoMLTablesListColumnSpecsOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations."
- )
def execute(self, context: Context):
hook = CloudAutoMLHook(
@@ -824,6 +796,12 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.dataset.UpdateDatasetOperator",
+ category=AirflowProviderDeprecationWarning,
+ reason="Shutdown of legacy version of AutoML Tables on March 31, 2024.",
+)
class AutoMLTablesUpdateDatasetOperator(GoogleCloudBaseOperator):
"""
Updates a dataset.
@@ -892,12 +870,6 @@ def __init__(
self.retry = retry
self.gcp_conn_id = gcp_conn_id
self.impersonation_chain = impersonation_chain
- raise AirflowException(
- "Operator AutoMLTablesUpdateDatasetOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations. "
- "Please use UpdateDatasetOperator from Vertex AI instead."
- )
def execute(self, context: Context):
hook = CloudAutoMLHook(
@@ -924,6 +896,11 @@ def execute(self, context: Context):
return Dataset.to_dict(result)
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.model_service.GetModelOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLGetModelOperator(GoogleCloudBaseOperator):
"""
Get Google Cloud AutoML model.
@@ -1003,8 +980,6 @@ def execute(self, context: Context):
timeout=self.timeout,
metadata=self.metadata,
)
- if not hasattr(result, "translation_model_metadata"):
- _raise_exception_for_deprecated_operator(self.__class__.__name__, "GetModelOperator")
model = Model.to_dict(result)
project_id = self.project_id or hook.project_id
if project_id:
@@ -1018,6 +993,12 @@ def execute(self, context: Context):
return model
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.model_service.DeleteModelOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateDeleteModelOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLDeleteModelOperator(GoogleCloudBaseOperator):
"""
Delete Google Cloud AutoML model.
@@ -1088,7 +1069,7 @@ def execute(self, context: Context):
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
)
- model: Model = hook.get_model(
+ hook.get_model(
model_id=self.model_id,
location=self.location,
project_id=self.project_id,
@@ -1096,8 +1077,6 @@ def execute(self, context: Context):
timeout=self.timeout,
metadata=self.metadata,
)
- if not hasattr(model, "translation_model_metadata"):
- _raise_exception_for_deprecated_operator(self.__class__.__name__, "DeleteModelOperator")
operation = hook.delete_model(
model_id=self.model_id,
location=self.location,
@@ -1110,6 +1089,11 @@ def execute(self, context: Context):
self.log.info("Deletion is completed")
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.endpoint_service.DeployModelOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLDeployModelOperator(GoogleCloudBaseOperator):
"""
Deploys a model; if a model is already deployed, deploying it with the same parameters has no effect.
@@ -1187,13 +1171,6 @@ def __init__(
self.retry = retry
self.gcp_conn_id = gcp_conn_id
self.impersonation_chain = impersonation_chain
- raise AirflowException(
- "Operator AutoMLDeployModelOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML AutoML Natural Language, Vision, Video Intelligence "
- "on March 31, 2024. "
- "For additional information see: https://cloud.google.com/vision/automl/docs/deprecations. "
- "Please use DeployModelOperator from Vertex AI instead."
- )
def execute(self, context: Context):
hook = CloudAutoMLHook(
@@ -1214,6 +1191,11 @@ def execute(self, context: Context):
self.log.info("Model was deployed successfully.")
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ category=AirflowProviderDeprecationWarning,
+ reason="Shutdown of legacy version of AutoML Tables on March 31, 2024.",
+)
class AutoMLTablesListTableSpecsOperator(GoogleCloudBaseOperator):
"""
Lists table specs in a dataset.
@@ -1288,11 +1270,6 @@ def __init__(
self.retry = retry
self.gcp_conn_id = gcp_conn_id
self.impersonation_chain = impersonation_chain
- raise AirflowException(
- "Operator AutoMLTablesListTableSpecsOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations. "
- )
def execute(self, context: Context):
hook = CloudAutoMLHook(
@@ -1324,6 +1301,12 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.dataset.ListDatasetsOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateDatasetsListOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLListDatasetOperator(GoogleCloudBaseOperator):
"""
Lists AutoML Datasets in project.
@@ -1399,14 +1382,7 @@ def execute(self, context: Context):
)
result = []
for dataset in page_iterator:
- if not hasattr(dataset, "translation_dataset_metadata"):
- warnings.warn(
- "Class `AutoMLListDatasetOperator` has been deprecated and no longer available. "
- "Please use `ListDatasetsOperator` instead.",
- stacklevel=2,
- )
- else:
- result.append(Dataset.to_dict(dataset))
+ result.append(Dataset.to_dict(dataset))
self.log.info("Datasets obtained.")
self.xcom_push(
@@ -1420,6 +1396,12 @@ def execute(self, context: Context):
return result
+@deprecated(
+ planned_removal_date="September 30, 2025",
+ use_instead="airflow.providers.google.cloud.operators.vertex_ai.dataset.ListDatasetsOperator, "
+ "airflow.providers.google.cloud.operators.translate.TranslateDatasetsListOperator",
+ category=AirflowProviderDeprecationWarning,
+)
class AutoMLDeleteDatasetOperator(GoogleCloudBaseOperator):
"""
Deletes a dataset and all of its contents.
@@ -1498,7 +1480,7 @@ def execute(self, context: Context):
gcp_conn_id=self.gcp_conn_id,
impersonation_chain=self.impersonation_chain,
)
- dataset: Dataset = hook.get_dataset(
+ hook.get_dataset(
dataset_id=self.dataset_id,
location=self.location,
project_id=self.project_id,
@@ -1506,8 +1488,6 @@ def execute(self, context: Context):
timeout=self.timeout,
metadata=self.metadata,
)
- if not hasattr(dataset, "translation_dataset_metadata"):
- _raise_exception_for_deprecated_operator(self.__class__.__name__, "DeleteDatasetOperator")
dataset_id_list = self._parse_dataset_id(self.dataset_id)
for dataset_id in dataset_id_list:
self.log.info("Deleting dataset %s", dataset_id)
diff --git a/providers/tests/google/cloud/operators/test_automl.py b/providers/tests/google/cloud/operators/test_automl.py
index e0b286a343257..94dca98be917b 100644
--- a/providers/tests/google/cloud/operators/test_automl.py
+++ b/providers/tests/google/cloud/operators/test_automl.py
@@ -28,7 +28,7 @@
from google.api_core.gapic_v1.method import DEFAULT
from google.cloud.automl_v1beta1 import BatchPredictResult, Dataset, Model, PredictResponse
-from airflow.exceptions import AirflowException, AirflowProviderDeprecationWarning
+from airflow.exceptions import AirflowProviderDeprecationWarning
from airflow.providers.google.cloud.hooks.automl import CloudAutoMLHook
from airflow.providers.google.cloud.hooks.vertex_ai.prediction_service import PredictionServiceHook
from airflow.providers.google.cloud.operators.automl import (
@@ -85,12 +85,14 @@ def test_execute(self, mock_hook):
mock_hook.return_value.create_model.return_value.result.return_value = Model(name=MODEL_PATH)
mock_hook.return_value.extract_object_id = extract_object_id
mock_hook.return_value.wait_for_operation.return_value = Model()
- op = AutoMLTrainModelOperator(
- model=MODEL,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLTrainModelOperator(
+ model=MODEL,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ task_id=TASK_ID,
+ )
+
op.execute(context=mock.MagicMock())
mock_hook.return_value.create_model.assert_called_once_with(
model=MODEL,
@@ -101,34 +103,19 @@ def test_execute(self, mock_hook):
metadata=(),
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- op = AutoMLTrainModelOperator(
- model=MODEL_DEPRECATED,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLTrainModelOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLTrainModelOperator,
- # Templated fields
- model="{{ 'model' }}",
- location="{{ 'location' }}",
- impersonation_chain="{{ 'impersonation_chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLTrainModelOperator,
+ # Templated fields
+ model="{{ 'model' }}",
+ location="{{ 'location' }}",
+ impersonation_chain="{{ 'impersonation_chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -177,38 +164,6 @@ def test_execute(self, mock_hook, mock_link_persist):
dataset_id=DATASET_ID,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- returned_model = mock.MagicMock()
- del returned_model.translation_model_metadata
- mock_hook.return_value.get_model.return_value = returned_model
- mock_hook.return_value.extract_object_id = extract_object_id
- with pytest.warns(AirflowProviderDeprecationWarning):
- op = AutoMLBatchPredictOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- input_config=INPUT_CONFIG,
- output_config=OUTPUT_CONFIG,
- task_id=TASK_ID,
- prediction_params={},
- )
- expected_exception_str = (
- "AutoMLBatchPredictOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.get_model.assert_called_once_with(
- location=GCP_LOCATION,
- model_id=MODEL_ID,
- project_id=GCP_PROJECT_ID,
- retry=DEFAULT,
- timeout=None,
- metadata=(),
- )
- mock_hook.return_value.batch_predict.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
with pytest.warns(AirflowProviderDeprecationWarning):
@@ -244,14 +199,15 @@ def test_execute(self, mock_hook, mock_link_persist):
mock_hook.return_value.predict.return_value = PredictResponse()
mock_hook.return_value.get_model.return_value = mock.MagicMock(**MODEL)
mock_context = {"ti": mock.MagicMock()}
- op = AutoMLPredictOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- payload=PAYLOAD,
- task_id=TASK_ID,
- operation_params={"TEST_KEY": "TEST_VALUE"},
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLPredictOperator(
+ model_id=MODEL_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ payload=PAYLOAD,
+ task_id=TASK_ID,
+ operation_params={"TEST_KEY": "TEST_VALUE"},
+ )
op.execute(context=mock_context)
mock_hook.return_value.predict.assert_called_once_with(
location=GCP_LOCATION,
@@ -273,18 +229,19 @@ def test_execute(self, mock_hook, mock_link_persist):
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLPredictOperator,
- # Templated fields
- model_id="{{ 'model-id' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- payload={},
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLPredictOperator,
+ # Templated fields
+ model_id="{{ 'model-id' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ payload={},
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -294,51 +251,27 @@ def test_templating(self, create_task_instance_of_operator, session):
assert task.location == "location"
assert task.impersonation_chain == "impersonation-chain"
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecation(self, mock_hook):
- returned_model = mock.MagicMock(**MODEL_DEPRECATED)
- del returned_model.translation_model_metadata
- mock_hook.return_value.get_model.return_value = returned_model
-
- mock_hook.return_value.predict.return_value = PredictResponse()
-
- op = AutoMLPredictOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- payload=PAYLOAD,
- task_id=TASK_ID,
- operation_params={"TEST_KEY": "TEST_VALUE"},
- )
- expected_exception_str = (
- "AutoMLPredictOperator for text, image, and video prediction has been "
- "deprecated. Please use endpoint_id param instead of model_id param."
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.predict.assert_not_called()
-
@pytest.mark.db_test
def test_hook_type(self):
- op = AutoMLPredictOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- payload=PAYLOAD,
- task_id=TASK_ID,
- operation_params={"TEST_KEY": "TEST_VALUE"},
- )
- assert isinstance(op.hook, CloudAutoMLHook)
-
- op = AutoMLPredictOperator(
- endpoint_id="endpoint_id",
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- payload=PAYLOAD,
- task_id=TASK_ID,
- operation_params={"TEST_KEY": "TEST_VALUE"},
- )
- assert isinstance(op.hook, PredictionServiceHook)
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLPredictOperator(
+ model_id=MODEL_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ payload=PAYLOAD,
+ task_id=TASK_ID,
+ operation_params={"TEST_KEY": "TEST_VALUE"},
+ )
+ assert isinstance(op.hook, CloudAutoMLHook)
+ op = AutoMLPredictOperator(
+ endpoint_id="endpoint_id",
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ payload=PAYLOAD,
+ task_id=TASK_ID,
+ operation_params={"TEST_KEY": "TEST_VALUE"},
+ )
+ assert isinstance(op.hook, PredictionServiceHook)
class TestAutoMLCreateImportOperator:
@@ -346,13 +279,13 @@ class TestAutoMLCreateImportOperator:
def test_execute(self, mock_hook):
mock_hook.return_value.create_dataset.return_value = Dataset(name=DATASET_PATH)
mock_hook.return_value.extract_object_id = extract_object_id
-
- op = AutoMLCreateDatasetOperator(
- dataset=DATASET,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLCreateDatasetOperator(
+ dataset=DATASET,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ task_id=TASK_ID,
+ )
op.execute(context=mock.MagicMock())
mock_hook.return_value.create_dataset.assert_called_once_with(
dataset=DATASET,
@@ -363,35 +296,20 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- op = AutoMLCreateDatasetOperator(
- dataset=DATASET_DEPRECATED,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLCreateDatasetOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.create_dataset.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLCreateDatasetOperator,
- # Templated fields
- dataset="{{ 'dataset' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLCreateDatasetOperator,
+ # Templated fields
+ dataset="{{ 'dataset' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -403,20 +321,14 @@ def test_templating(self, create_task_instance_of_operator, session):
class TestAutoMLTablesListColumnsSpecsOperator:
- expected_exception_string = (
- "Operator AutoMLTablesListColumnSpecsOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations."
- )
-
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
table_spec = "table_spec_id"
filter_ = "filter"
page_size = 42
- with pytest.raises(AirflowException, match=self.expected_exception_string):
- _ = AutoMLTablesListColumnSpecsOperator(
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ AutoMLTablesListColumnSpecsOperator(
dataset_id=DATASET_ID,
table_spec_id=table_spec,
location=GCP_LOCATION,
@@ -430,8 +342,8 @@ def test_execute(self, mock_hook):
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator):
- with pytest.raises(AirflowException, match=self.expected_exception_string):
- _ = create_task_instance_of_operator(
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ create_task_instance_of_operator(
AutoMLTablesListColumnSpecsOperator,
# Templated fields
dataset_id="{{ 'dataset-id' }}",
@@ -448,21 +360,13 @@ def test_templating(self, create_task_instance_of_operator):
class TestAutoMLTablesUpdateDatasetOperator:
- expected_exception_string = (
- "Operator AutoMLTablesUpdateDatasetOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations. "
- "Please use UpdateDatasetOperator from Vertex AI instead."
- )
-
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
mock_hook.return_value.update_dataset.return_value = Dataset(name=DATASET_PATH)
-
dataset = copy.deepcopy(DATASET)
dataset["name"] = DATASET_ID
- with pytest.raises(AirflowException, match=self.expected_exception_string):
+ with pytest.warns(AirflowProviderDeprecationWarning):
AutoMLTablesUpdateDatasetOperator(
dataset=dataset,
update_mask=MASK,
@@ -473,7 +377,7 @@ def test_execute(self, mock_hook):
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator):
- with pytest.raises(AirflowException, match=self.expected_exception_string):
+ with pytest.warns(AirflowProviderDeprecationWarning):
create_task_instance_of_operator(
AutoMLTablesUpdateDatasetOperator,
# Templated fields
@@ -492,13 +396,13 @@ class TestAutoMLGetModelOperator:
def test_execute(self, mock_hook):
mock_hook.return_value.get_model.return_value = Model(name=MODEL_PATH)
mock_hook.return_value.extract_object_id = extract_object_id
-
- op = AutoMLGetModelOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLGetModelOperator(
+ model_id=MODEL_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ task_id=TASK_ID,
+ )
op.execute(context=mock.MagicMock())
mock_hook.return_value.get_model.assert_called_once_with(
location=GCP_LOCATION,
@@ -509,46 +413,20 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- returned_model = mock.MagicMock(**MODEL_DEPRECATED)
- del returned_model.translation_model_metadata
- mock_hook.return_value.get_model.return_value = returned_model
-
- op = AutoMLGetModelOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLGetModelOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.get_model.assert_called_once_with(
- location=GCP_LOCATION,
- metadata=(),
- model_id=MODEL_ID,
- project_id=GCP_PROJECT_ID,
- retry=DEFAULT,
- timeout=None,
- )
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLGetModelOperator,
- # Templated fields
- model_id="{{ 'model-id' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLGetModelOperator,
+ # Templated fields
+ model_id="{{ 'model-id' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -562,12 +440,13 @@ def test_templating(self, create_task_instance_of_operator, session):
class TestAutoMLDeleteModelOperator:
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
- op = AutoMLDeleteModelOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLDeleteModelOperator(
+ model_id=MODEL_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ task_id=TASK_ID,
+ )
op.execute(context=None)
mock_hook.return_value.delete_model.assert_called_once_with(
location=GCP_LOCATION,
@@ -578,47 +457,20 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- returned_model = mock.MagicMock(**MODEL_DEPRECATED)
- del returned_model.translation_model_metadata
- mock_hook.return_value.get_model.return_value = returned_model
-
- op = AutoMLDeleteModelOperator(
- model_id=MODEL_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLDeleteModelOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.get_model.assert_called_once_with(
- location=GCP_LOCATION,
- metadata=(),
- model_id=MODEL_ID,
- project_id=GCP_PROJECT_ID,
- retry=DEFAULT,
- timeout=None,
- )
- mock_hook.return_value.delete_model.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLDeleteModelOperator,
- # Templated fields
- model_id="{{ 'model-id' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLDeleteModelOperator,
+ # Templated fields
+ model_id="{{ 'model-id' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -630,19 +482,11 @@ def test_templating(self, create_task_instance_of_operator, session):
class TestAutoMLDeployModelOperator:
- expected_exception_string = (
- "Operator AutoMLDeployModelOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML AutoML Natural Language, Vision, Video Intelligence "
- "on March 31, 2024. "
- "For additional information see: https://cloud.google.com/vision/automl/docs/deprecations. "
- "Please use DeployModelOperator from Vertex AI instead."
- )
-
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
image_detection_metadata = {}
- with pytest.raises(AirflowException, match=self.expected_exception_string):
+ with pytest.warns(AirflowProviderDeprecationWarning):
AutoMLDeployModelOperator(
model_id=MODEL_ID,
image_detection_metadata=image_detection_metadata,
@@ -655,7 +499,7 @@ def test_execute(self, mock_hook):
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator):
- with pytest.raises(AirflowException, match=self.expected_exception_string):
+ with pytest.warns(AirflowProviderDeprecationWarning):
create_task_instance_of_operator(
AutoMLDeployModelOperator,
# Templated fields
@@ -672,13 +516,14 @@ def test_templating(self, create_task_instance_of_operator):
class TestAutoMLDatasetImportOperator:
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
- op = AutoMLImportDataOperator(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- input_config=INPUT_CONFIG,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLImportDataOperator(
+ dataset_id=DATASET_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ input_config=INPUT_CONFIG,
+ task_id=TASK_ID,
+ )
op.execute(context=mock.MagicMock())
mock_hook.return_value.import_data.assert_called_once_with(
input_config=INPUT_CONFIG,
@@ -690,49 +535,21 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- returned_dataset = mock.MagicMock()
- del returned_dataset.translation_dataset_metadata
- mock_hook.return_value.get_dataset.return_value = returned_dataset
-
- op = AutoMLImportDataOperator(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- input_config=INPUT_CONFIG,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLImportDataOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.get_dataset.assert_called_once_with(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- metadata=(),
- retry=DEFAULT,
- timeout=None,
- )
- mock_hook.return_value.import_data.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLImportDataOperator,
- # Templated fields
- dataset_id="{{ 'dataset-id' }}",
- input_config="{{ 'input-config' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLImportDataOperator,
+ # Templated fields
+ dataset_id="{{ 'dataset-id' }}",
+ input_config="{{ 'input-config' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -745,19 +562,13 @@ def test_templating(self, create_task_instance_of_operator, session):
class TestAutoMLTablesListTableSpecsOperator:
- expected_exception_string = (
- "Operator AutoMLTablesListTableSpecsOperator has been deprecated due to shutdown of "
- "a legacy version of AutoML Tables on March 31, 2024. "
- "For additional information see: https://cloud.google.com/automl-tables/docs/deprecations. "
- )
-
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
filter_ = "filter"
page_size = 42
- with pytest.raises(AirflowException, match=self.expected_exception_string):
- _ = AutoMLTablesListTableSpecsOperator(
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ AutoMLTablesListTableSpecsOperator(
dataset_id=DATASET_ID,
location=GCP_LOCATION,
project_id=GCP_PROJECT_ID,
@@ -769,8 +580,8 @@ def test_execute(self, mock_hook):
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator):
- with pytest.raises(AirflowException, match=self.expected_exception_string):
- _ = create_task_instance_of_operator(
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ create_task_instance_of_operator(
AutoMLTablesListTableSpecsOperator,
# Templated fields
dataset_id="{{ 'dataset-id' }}",
@@ -787,7 +598,8 @@ def test_templating(self, create_task_instance_of_operator):
class TestAutoMLDatasetListOperator:
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
- op = AutoMLListDatasetOperator(location=GCP_LOCATION, project_id=GCP_PROJECT_ID, task_id=TASK_ID)
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLListDatasetOperator(location=GCP_LOCATION, project_id=GCP_PROJECT_ID, task_id=TASK_ID)
op.execute(context=mock.MagicMock())
mock_hook.return_value.list_datasets.assert_called_once_with(
location=GCP_LOCATION,
@@ -797,39 +609,19 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- not_valid_dataset = mock.MagicMock()
- del not_valid_dataset.translation_dataset_metadata
- mock_hook.return_value.list_datasets.return_value = [DATASET, not_valid_dataset]
- op = AutoMLListDatasetOperator(location=GCP_LOCATION, project_id=GCP_PROJECT_ID, task_id=TASK_ID)
- expected_warning_str = (
- "Class `AutoMLListDatasetOperator` has been deprecated and no longer available. "
- "Please use `ListDatasetsOperator` instead"
- )
- with pytest.warns(UserWarning, match=expected_warning_str):
- op.execute(context=mock.MagicMock())
-
- mock_hook.return_value.list_datasets.assert_called_once_with(
- location=GCP_LOCATION,
- metadata=(),
- project_id=GCP_PROJECT_ID,
- retry=DEFAULT,
- timeout=None,
- )
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLListDatasetOperator,
- # Templated fields
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLListDatasetOperator,
+ # Templated fields
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
@@ -842,12 +634,13 @@ def test_templating(self, create_task_instance_of_operator, session):
class TestAutoMLDatasetDeleteOperator:
@mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
def test_execute(self, mock_hook):
- op = AutoMLDeleteDatasetOperator(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ op = AutoMLDeleteDatasetOperator(
+ dataset_id=DATASET_ID,
+ location=GCP_LOCATION,
+ project_id=GCP_PROJECT_ID,
+ task_id=TASK_ID,
+ )
op.execute(context=None)
mock_hook.return_value.delete_dataset.assert_called_once_with(
location=GCP_LOCATION,
@@ -858,47 +651,20 @@ def test_execute(self, mock_hook):
timeout=None,
)
- @mock.patch("airflow.providers.google.cloud.operators.automl.CloudAutoMLHook")
- def test_execute_deprecated(self, mock_hook):
- returned_dataset = mock.MagicMock()
- del returned_dataset.translation_dataset_metadata
- mock_hook.return_value.get_dataset.return_value = returned_dataset
-
- op = AutoMLDeleteDatasetOperator(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- task_id=TASK_ID,
- )
- expected_exception_str = (
- "AutoMLDeleteDatasetOperator for text, image, and video prediction has been "
- "deprecated and no longer available"
- )
- with pytest.raises(AirflowException, match=expected_exception_str):
- op.execute(context=mock.MagicMock())
- mock_hook.return_value.get_dataset.assert_called_once_with(
- dataset_id=DATASET_ID,
- location=GCP_LOCATION,
- project_id=GCP_PROJECT_ID,
- metadata=(),
- retry=DEFAULT,
- timeout=None,
- )
- mock_hook.return_value.delete_dataset.assert_not_called()
-
@pytest.mark.db_test
def test_templating(self, create_task_instance_of_operator, session):
- ti = create_task_instance_of_operator(
- AutoMLDeleteDatasetOperator,
- # Templated fields
- dataset_id="{{ 'dataset-id' }}",
- location="{{ 'location' }}",
- project_id="{{ 'project-id' }}",
- impersonation_chain="{{ 'impersonation-chain' }}",
- # Other parameters
- dag_id="test_template_body_templating_dag",
- task_id="test_template_body_templating_task",
- )
+ with pytest.warns(AirflowProviderDeprecationWarning):
+ ti = create_task_instance_of_operator(
+ AutoMLDeleteDatasetOperator,
+ # Templated fields
+ dataset_id="{{ 'dataset-id' }}",
+ location="{{ 'location' }}",
+ project_id="{{ 'project-id' }}",
+ impersonation_chain="{{ 'impersonation-chain' }}",
+ # Other parameters
+ dag_id="test_template_body_templating_dag",
+ task_id="test_template_body_templating_task",
+ )
session.add(ti)
session.commit()
ti.render_templates()
diff --git a/providers/tests/system/google/cloud/automl/__init__.py b/providers/tests/system/google/cloud/automl/__init__.py
deleted file mode 100644
index 13a83393a9124..0000000000000
--- a/providers/tests/system/google/cloud/automl/__init__.py
+++ /dev/null
@@ -1,16 +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.
diff --git a/providers/tests/system/google/cloud/automl/example_automl_dataset.py b/providers/tests/system/google/cloud/automl/example_automl_dataset.py
deleted file mode 100644
index d1305da4fb081..0000000000000
--- a/providers/tests/system/google/cloud/automl/example_automl_dataset.py
+++ /dev/null
@@ -1,175 +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 for Google AutoML service testing dataset operations."""
-
-from __future__ import annotations
-
-import os
-from datetime import datetime
-
-from google.cloud import storage # type: ignore[attr-defined]
-
-from airflow.decorators import task
-from airflow.models.dag import DAG
-from airflow.providers.google.cloud.operators.automl import (
- AutoMLCreateDatasetOperator,
- AutoMLDeleteDatasetOperator,
- AutoMLImportDataOperator,
- AutoMLListDatasetOperator,
-)
-from airflow.providers.google.cloud.operators.gcs import (
- GCSCreateBucketOperator,
- GCSDeleteBucketOperator,
-)
-from airflow.providers.google.cloud.transfers.gcs_to_gcs import GCSToGCSOperator
-from airflow.utils.trigger_rule import TriggerRule
-
-ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
-DAG_ID = "automl_dataset"
-GCP_PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT", "default")
-
-GCP_AUTOML_LOCATION = "us-central1"
-RESOURCE_DATA_BUCKET = "airflow-system-tests-resources"
-DATA_SAMPLE_GCS_BUCKET_NAME = f"bucket_{DAG_ID}_{ENV_ID}".replace("_", "-")
-
-DATASET_NAME = f"ds_{DAG_ID}_{ENV_ID}".replace("-", "_")
-DATASET = {
- "display_name": DATASET_NAME,
- "translation_dataset_metadata": {
- "source_language_code": "en",
- "target_language_code": "es",
- },
-}
-
-CSV_FILE_NAME = "en-es.csv"
-TSV_FILE_NAME = "en-es.tsv"
-GCS_FILE_PATH = f"automl/datasets/translate/{CSV_FILE_NAME}"
-AUTOML_DATASET_BUCKET = f"gs://{DATA_SAMPLE_GCS_BUCKET_NAME}/automl/{CSV_FILE_NAME}"
-IMPORT_INPUT_CONFIG = {"gcs_source": {"input_uris": [AUTOML_DATASET_BUCKET]}}
-
-
-with DAG(
- dag_id=DAG_ID,
- schedule="@once",
- start_date=datetime(2021, 1, 1),
- catchup=False,
- tags=["example", "automl", "dataset"],
-) as dag:
- create_bucket = GCSCreateBucketOperator(
- task_id="create_bucket",
- bucket_name=DATA_SAMPLE_GCS_BUCKET_NAME,
- storage_class="REGIONAL",
- location=GCP_AUTOML_LOCATION,
- )
-
- @task
- def upload_updated_csv_file_to_gcs():
- # download file into memory
- storage_client = storage.Client()
- bucket = storage_client.bucket(RESOURCE_DATA_BUCKET, GCP_PROJECT_ID)
- blob = bucket.blob(GCS_FILE_PATH)
- contents = blob.download_as_string().decode()
-
- # update file content
- updated_contents = contents.replace("template-bucket", DATA_SAMPLE_GCS_BUCKET_NAME)
-
- # upload updated content to bucket
- destination_bucket = storage_client.bucket(DATA_SAMPLE_GCS_BUCKET_NAME)
- destination_blob = destination_bucket.blob(f"automl/{CSV_FILE_NAME}")
- destination_blob.upload_from_string(updated_contents)
-
- # AutoML requires a .csv file with links to .tsv/.tmx files containing translation training data
- upload_csv_dataset_file = upload_updated_csv_file_to_gcs()
-
- # The .tsv file contains training data with translated language pairs
- copy_tsv_dataset_file = GCSToGCSOperator(
- task_id="copy_dataset_file",
- source_bucket=RESOURCE_DATA_BUCKET,
- source_object=f"automl/datasets/translate/{TSV_FILE_NAME}",
- destination_bucket=DATA_SAMPLE_GCS_BUCKET_NAME,
- destination_object=f"automl/{TSV_FILE_NAME}",
- )
-
- # [START howto_operator_automl_create_dataset]
- create_dataset = AutoMLCreateDatasetOperator(
- task_id="create_dataset",
- dataset=DATASET,
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
- dataset_id = create_dataset.output["dataset_id"]
- # [END howto_operator_automl_create_dataset]
-
- # [START howto_operator_automl_import_data]
- import_dataset = AutoMLImportDataOperator(
- task_id="import_dataset",
- dataset_id=dataset_id,
- location=GCP_AUTOML_LOCATION,
- input_config=IMPORT_INPUT_CONFIG,
- )
- # [END howto_operator_automl_import_data]
-
- # [START howto_operator_list_dataset]
- list_datasets = AutoMLListDatasetOperator(
- task_id="list_datasets",
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
- # [END howto_operator_list_dataset]
-
- # [START howto_operator_delete_dataset]
- delete_dataset = AutoMLDeleteDatasetOperator(
- task_id="delete_dataset",
- dataset_id=dataset_id,
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- trigger_rule=TriggerRule.ALL_DONE,
- )
- # [END howto_operator_delete_dataset]
-
- delete_bucket = GCSDeleteBucketOperator(
- task_id="delete_bucket",
- bucket_name=DATA_SAMPLE_GCS_BUCKET_NAME,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- (
- # TEST SETUP
- [create_bucket >> upload_csv_dataset_file >> copy_tsv_dataset_file]
- # create_bucket
- >> create_dataset
- # TEST BODY
- >> import_dataset
- >> list_datasets
- # TEST TEARDOWN
- >> delete_dataset
- >> 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: tests/system/README.md#run_via_pytest)
-test_run = get_test_run(dag)
diff --git a/providers/tests/system/google/cloud/automl/example_automl_translation.py b/providers/tests/system/google/cloud/automl/example_automl_translation.py
deleted file mode 100644
index e758bdf113e48..0000000000000
--- a/providers/tests/system/google/cloud/automl/example_automl_translation.py
+++ /dev/null
@@ -1,201 +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 that uses Google AutoML Translation services."""
-
-from __future__ import annotations
-
-import os
-from datetime import datetime
-from typing import cast
-
-# The storage module cannot be imported yet https://github.com/googleapis/python-storage/issues/393
-from google.cloud import storage # type: ignore[attr-defined]
-
-from airflow.decorators import task
-from airflow.models.dag import DAG
-from airflow.models.xcom_arg import XComArg
-from airflow.providers.google.cloud.operators.automl import (
- AutoMLCreateDatasetOperator,
- AutoMLDeleteDatasetOperator,
- AutoMLDeleteModelOperator,
- AutoMLGetModelOperator,
- AutoMLImportDataOperator,
- AutoMLPredictOperator,
- AutoMLTrainModelOperator,
-)
-from airflow.providers.google.cloud.operators.gcs import GCSCreateBucketOperator, GCSDeleteBucketOperator
-from airflow.providers.google.cloud.transfers.gcs_to_gcs import GCSToGCSOperator
-from airflow.utils.trigger_rule import TriggerRule
-
-DAG_ID = "automl_translate"
-GCP_PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT", "default")
-ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
-GCP_AUTOML_LOCATION = "us-central1"
-DATA_SAMPLE_GCS_BUCKET_NAME = f"bucket_{DAG_ID}_{ENV_ID}".replace("_", "-")
-RESOURCE_DATA_BUCKET = "airflow-system-tests-resources"
-
-
-MODEL_NAME = "translate_test_model"
-MODEL = {
- "display_name": MODEL_NAME,
- "translation_model_metadata": {},
-}
-
-DATASET_NAME = f"ds_{DAG_ID}_{ENV_ID}".replace("-", "_")
-DATASET = {
- "display_name": DATASET_NAME,
- "translation_dataset_metadata": {
- "source_language_code": "en",
- "target_language_code": "es",
- },
-}
-
-CSV_FILE_NAME = "en-es.csv"
-TSV_FILE_NAME = "en-es.tsv"
-GCS_FILE_PATH = f"automl/datasets/translate/{CSV_FILE_NAME}"
-AUTOML_DATASET_BUCKET = f"gs://{DATA_SAMPLE_GCS_BUCKET_NAME}/automl/{CSV_FILE_NAME}"
-IMPORT_INPUT_CONFIG = {"gcs_source": {"input_uris": [AUTOML_DATASET_BUCKET]}}
-
-
-# Example DAG for AutoML Translation
-with DAG(
- DAG_ID,
- schedule="@once",
- start_date=datetime(2021, 1, 1),
- catchup=False,
- tags=["example", "automl", "translate"],
-) as dag:
- create_bucket = GCSCreateBucketOperator(
- task_id="create_bucket",
- bucket_name=DATA_SAMPLE_GCS_BUCKET_NAME,
- storage_class="REGIONAL",
- location=GCP_AUTOML_LOCATION,
- )
-
- @task
- def upload_csv_file_to_gcs():
- # download file into memory
- storage_client = storage.Client()
- bucket = storage_client.bucket(RESOURCE_DATA_BUCKET)
- blob = bucket.blob(GCS_FILE_PATH)
- contents = blob.download_as_string().decode()
-
- # update memory content
- updated_contents = contents.replace("template-bucket", DATA_SAMPLE_GCS_BUCKET_NAME)
-
- # upload updated content to bucket
- destination_bucket = storage_client.bucket(DATA_SAMPLE_GCS_BUCKET_NAME)
- destination_blob = destination_bucket.blob(f"automl/{CSV_FILE_NAME}")
- destination_blob.upload_from_string(updated_contents)
-
- upload_csv_file_to_gcs_task = upload_csv_file_to_gcs()
-
- copy_dataset_file = GCSToGCSOperator(
- task_id="copy_dataset_file",
- source_bucket=RESOURCE_DATA_BUCKET,
- source_object=f"automl/datasets/translate/{TSV_FILE_NAME}",
- destination_bucket=DATA_SAMPLE_GCS_BUCKET_NAME,
- destination_object=f"automl/{TSV_FILE_NAME}",
- )
-
- create_dataset = AutoMLCreateDatasetOperator(
- task_id="create_dataset", dataset=DATASET, location=GCP_AUTOML_LOCATION
- )
-
- dataset_id = cast(str, XComArg(create_dataset, key="dataset_id"))
-
- import_dataset = AutoMLImportDataOperator(
- task_id="import_dataset",
- dataset_id=dataset_id,
- location=GCP_AUTOML_LOCATION,
- input_config=IMPORT_INPUT_CONFIG,
- )
-
- MODEL["dataset_id"] = dataset_id
- # [START howto_operator_automl_create_model]
- create_model = AutoMLTrainModelOperator(task_id="create_model", model=MODEL, location=GCP_AUTOML_LOCATION)
- # [END howto_operator_automl_create_model]
- model_id = cast(str, XComArg(create_model, key="model_id"))
-
- # [START howto_operator_get_model]
- get_model = AutoMLGetModelOperator(
- task_id="get_model",
- model_id=model_id,
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
- # [END howto_operator_get_model]
-
- # [START howto_operator_prediction]
- TRANSLATION_STR = "A Dog walks down the street"
- predict_task = AutoMLPredictOperator(
- task_id="predict_task",
- model_id=model_id,
- payload={"text_snippet": {"content": TRANSLATION_STR}},
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
- # [END howto_operator_prediction]
-
- delete_model = AutoMLDeleteModelOperator(
- task_id="delete_model",
- model_id=model_id,
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
-
- delete_dataset = AutoMLDeleteDatasetOperator(
- task_id="delete_dataset",
- dataset_id=dataset_id,
- location=GCP_AUTOML_LOCATION,
- project_id=GCP_PROJECT_ID,
- )
-
- delete_bucket = GCSDeleteBucketOperator(
- task_id="delete_bucket",
- bucket_name=DATA_SAMPLE_GCS_BUCKET_NAME,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- (
- # TEST SETUP
- [create_bucket >> upload_csv_file_to_gcs_task >> copy_dataset_file]
- # TEST BODY
- >> create_dataset
- >> import_dataset
- >> create_model
- >> get_model
- >> predict_task
- # TEST TEARDOWN
- >> delete_dataset
- >> delete_model
- >> 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: tests/system/README.md#run_via_pytest)
-test_run = get_test_run(dag)
diff --git a/providers/tests/system/google/cloud/automl/example_automl_video_classification.py b/providers/tests/system/google/cloud/automl/example_automl_video_classification.py
deleted file mode 100644
index 538831d307ca5..0000000000000
--- a/providers/tests/system/google/cloud/automl/example_automl_video_classification.py
+++ /dev/null
@@ -1,170 +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 that uses Google AutoML services.
-"""
-
-from __future__ import annotations
-
-import os
-from datetime import datetime
-
-from google.cloud.aiplatform import schema
-from google.protobuf.struct_pb2 import Value
-
-from airflow.models.dag import DAG
-from airflow.providers.google.cloud.operators.gcs import (
- GCSCreateBucketOperator,
- GCSDeleteBucketOperator,
- GCSSynchronizeBucketsOperator,
-)
-from airflow.providers.google.cloud.operators.vertex_ai.auto_ml import (
- CreateAutoMLVideoTrainingJobOperator,
- DeleteAutoMLTrainingJobOperator,
-)
-from airflow.providers.google.cloud.operators.vertex_ai.dataset import (
- CreateDatasetOperator,
- DeleteDatasetOperator,
- ImportDataOperator,
-)
-from airflow.utils.trigger_rule import TriggerRule
-
-ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
-PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT", "default")
-DAG_ID = "automl_video_clss"
-REGION = "us-central1"
-VIDEO_DISPLAY_NAME = f"auto-ml-video-clss-{ENV_ID}"
-MODEL_DISPLAY_NAME = f"auto-ml-video-clss-model-{ENV_ID}"
-
-RESOURCE_DATA_BUCKET = "airflow-system-tests-resources"
-VIDEO_GCS_BUCKET_NAME = f"bucket_video_clss_{ENV_ID}".replace("_", "-")
-
-VIDEO_DATASET = {
- "display_name": f"video-dataset-{ENV_ID}",
- "metadata_schema_uri": schema.dataset.metadata.video,
- "metadata": Value(string_value="video-dataset"),
-}
-VIDEO_DATA_CONFIG = [
- {
- "import_schema_uri": schema.dataset.ioformat.video.classification,
- "gcs_source": {"uris": [f"gs://{VIDEO_GCS_BUCKET_NAME}/automl/classification.csv"]},
- },
-]
-
-
-# Example DAG for AutoML Video Intelligence Classification
-with DAG(
- DAG_ID,
- schedule="@once",
- start_date=datetime(2021, 1, 1),
- catchup=False,
- tags=["example", "automl", "video", "classification"],
-) as dag:
- create_bucket = GCSCreateBucketOperator(
- task_id="create_bucket",
- bucket_name=VIDEO_GCS_BUCKET_NAME,
- storage_class="REGIONAL",
- location=REGION,
- )
-
- move_dataset_file = GCSSynchronizeBucketsOperator(
- task_id="move_dataset_to_bucket",
- source_bucket=RESOURCE_DATA_BUCKET,
- source_object="automl/datasets/video",
- destination_bucket=VIDEO_GCS_BUCKET_NAME,
- destination_object="automl",
- recursive=True,
- )
-
- create_video_dataset = CreateDatasetOperator(
- task_id="video_dataset",
- dataset=VIDEO_DATASET,
- region=REGION,
- project_id=PROJECT_ID,
- )
- video_dataset_id = create_video_dataset.output["dataset_id"]
-
- import_video_dataset = ImportDataOperator(
- task_id="import_video_data",
- dataset_id=video_dataset_id,
- region=REGION,
- project_id=PROJECT_ID,
- import_configs=VIDEO_DATA_CONFIG,
- )
-
- # [START howto_cloud_create_video_classification_training_job_operator]
- create_auto_ml_video_training_job = CreateAutoMLVideoTrainingJobOperator(
- task_id="auto_ml_video_task",
- display_name=VIDEO_DISPLAY_NAME,
- prediction_type="classification",
- model_type="CLOUD",
- dataset_id=video_dataset_id,
- model_display_name=MODEL_DISPLAY_NAME,
- region=REGION,
- project_id=PROJECT_ID,
- )
- # [END howto_cloud_create_video_classification_training_job_operator]
-
- delete_auto_ml_video_training_job = DeleteAutoMLTrainingJobOperator(
- task_id="delete_auto_ml_video_training_job",
- training_pipeline_id="{{ task_instance.xcom_pull(task_ids='auto_ml_video_task', "
- "key='training_id') }}",
- region=REGION,
- project_id=PROJECT_ID,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- delete_video_dataset = DeleteDatasetOperator(
- task_id="delete_video_dataset",
- dataset_id=video_dataset_id,
- region=REGION,
- project_id=PROJECT_ID,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- delete_bucket = GCSDeleteBucketOperator(
- task_id="delete_bucket",
- bucket_name=VIDEO_GCS_BUCKET_NAME,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- (
- # TEST SETUP
- [
- create_bucket >> move_dataset_file,
- create_video_dataset,
- ]
- >> import_video_dataset
- # TEST BODY
- >> create_auto_ml_video_training_job
- # TEST TEARDOWN
- >> delete_auto_ml_video_training_job
- >> delete_video_dataset
- >> 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: tests/system/README.md#run_via_pytest)
-test_run = get_test_run(dag)
diff --git a/providers/tests/system/google/cloud/automl/example_automl_vision_classification.py b/providers/tests/system/google/cloud/automl/example_automl_vision_classification.py
deleted file mode 100644
index e329febf5d9fc..0000000000000
--- a/providers/tests/system/google/cloud/automl/example_automl_vision_classification.py
+++ /dev/null
@@ -1,144 +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 that uses Google AutoML services.
-"""
-
-from __future__ import annotations
-
-import os
-from datetime import datetime
-
-from google.cloud.aiplatform import schema
-from google.protobuf.struct_pb2 import Value
-
-from airflow.models.dag import DAG
-from airflow.providers.google.cloud.operators.vertex_ai.auto_ml import (
- CreateAutoMLImageTrainingJobOperator,
- DeleteAutoMLTrainingJobOperator,
-)
-from airflow.providers.google.cloud.operators.vertex_ai.dataset import (
- CreateDatasetOperator,
- DeleteDatasetOperator,
- ImportDataOperator,
-)
-from airflow.utils.trigger_rule import TriggerRule
-
-DAG_ID = "automl_vision_clss"
-ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID", "default")
-PROJECT_ID = os.environ.get("SYSTEM_TESTS_GCP_PROJECT", "default")
-REGION = "us-central1"
-IMAGE_DISPLAY_NAME = f"automl-vision-clss-{ENV_ID}"
-MODEL_DISPLAY_NAME = f"automl-vision-clss-model-{ENV_ID}"
-
-RESOURCE_IMPORT_DATA_URI = (
- "gs://airflow-system-tests-resources/automl/datasets/vision/img_classification_short.csv"
-)
-IMAGE_DATASET = {
- "display_name": f"automl-vision-clss-dataset-{ENV_ID}",
- "metadata_schema_uri": schema.dataset.metadata.image,
- "metadata": Value(string_value="image-dataset"),
-}
-IMAGE_DATA_CONFIG = [
- {
- "import_schema_uri": schema.dataset.ioformat.image.single_label_classification,
- "gcs_source": {"uris": [RESOURCE_IMPORT_DATA_URI]},
- },
-]
-
-# Example DAG for AutoML Vision Classification
-with DAG(
- DAG_ID,
- schedule="@once", # Override to match your needs
- start_date=datetime(2021, 1, 1),
- catchup=False,
- tags=["example", "automl", "vision", "classification"],
-) as dag:
- create_image_dataset = CreateDatasetOperator(
- task_id="image_dataset",
- dataset=IMAGE_DATASET,
- region=REGION,
- project_id=PROJECT_ID,
- )
- image_dataset_id = create_image_dataset.output["dataset_id"]
-
- import_image_dataset = ImportDataOperator(
- task_id="import_image_data",
- dataset_id=image_dataset_id,
- region=REGION,
- project_id=PROJECT_ID,
- import_configs=IMAGE_DATA_CONFIG,
- )
-
- # [START howto_cloud_create_image_classification_training_job_operator]
- create_auto_ml_image_training_job = CreateAutoMLImageTrainingJobOperator(
- task_id="auto_ml_image_task",
- display_name=IMAGE_DISPLAY_NAME,
- dataset_id=image_dataset_id,
- prediction_type="classification",
- multi_label=False,
- model_type="CLOUD",
- training_fraction_split=0.6,
- validation_fraction_split=0.2,
- test_fraction_split=0.2,
- budget_milli_node_hours=8000,
- model_display_name=MODEL_DISPLAY_NAME,
- disable_early_stopping=False,
- region=REGION,
- project_id=PROJECT_ID,
- )
- # [END howto_cloud_create_image_classification_training_job_operator]
-
- delete_auto_ml_image_training_job = DeleteAutoMLTrainingJobOperator(
- task_id="delete_auto_ml_training_job",
- training_pipeline_id="{{ task_instance.xcom_pull(task_ids='auto_ml_image_task', "
- "key='training_id') }}",
- region=REGION,
- project_id=PROJECT_ID,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- delete_image_dataset = DeleteDatasetOperator(
- task_id="delete_image_dataset",
- dataset_id=image_dataset_id,
- region=REGION,
- project_id=PROJECT_ID,
- trigger_rule=TriggerRule.ALL_DONE,
- )
-
- (
- # TEST SETUP
- create_image_dataset
- >> import_image_dataset
- # TEST BODY
- >> create_auto_ml_image_training_job
- # TEST TEARDOWN
- >> delete_auto_ml_image_training_job
- >> delete_image_dataset
- )
-
- 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: tests/system/README.md#run_via_pytest)
-test_run = get_test_run(dag)
diff --git a/providers/tests/system/google/cloud/automl/resources/__init__.py b/providers/tests/system/google/cloud/automl/resources/__init__.py
deleted file mode 100644
index 13a83393a9124..0000000000000
--- a/providers/tests/system/google/cloud/automl/resources/__init__.py
+++ /dev/null
@@ -1,16 +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.
diff --git a/providers/tests/system/google/cloud/automl/example_automl_vision_object_detection.py b/providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_image_object_detection.py
similarity index 94%
rename from providers/tests/system/google/cloud/automl/example_automl_vision_object_detection.py
rename to providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_image_object_detection.py
index b87b42008d943..7a0796c6798c4 100644
--- a/providers/tests/system/google/cloud/automl/example_automl_vision_object_detection.py
+++ b/providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_image_object_detection.py
@@ -69,7 +69,7 @@
schedule="@once", # Override to match your needs
start_date=datetime(2021, 1, 1),
catchup=False,
- tags=["example", "automl", "vision", "object-detection"],
+ tags=["example", "vertex_ai", "automl", "vision", "object-detection"],
) as dag:
create_image_dataset = CreateDatasetOperator(
task_id="image_dataset",
@@ -86,8 +86,7 @@
project_id=PROJECT_ID,
import_configs=IMAGE_DATA_CONFIG,
)
-
- # [START howto_cloud_create_image_object_detection_training_job_operator]
+ # [START how_to_cloud_vertex_ai_create_auto_ml_image_object_detection_training_job_operator]
create_auto_ml_image_training_job = CreateAutoMLImageTrainingJobOperator(
task_id="auto_ml_image_task",
display_name=IMAGE_DISPLAY_NAME,
@@ -104,7 +103,7 @@
region=REGION,
project_id=PROJECT_ID,
)
- # [END howto_cloud_create_image_object_detection_training_job_operator]
+ # [END how_to_cloud_vertex_ai_create_auto_ml_image_object_detection_training_job_operator]
delete_auto_ml_image_training_job = DeleteAutoMLTrainingJobOperator(
task_id="delete_auto_ml_training_job",
diff --git a/providers/tests/system/google/cloud/automl/example_automl_video_tracking.py b/providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_video_tracking.py
similarity index 95%
rename from providers/tests/system/google/cloud/automl/example_automl_video_tracking.py
rename to providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_video_tracking.py
index 3876e9b0a39de..38d1e72024d4e 100644
--- a/providers/tests/system/google/cloud/automl/example_automl_video_tracking.py
+++ b/providers/tests/system/google/cloud/vertex_ai/example_vertex_ai_auto_ml_video_tracking.py
@@ -71,9 +71,9 @@
with DAG(
DAG_ID,
schedule="@once",
- start_date=datetime(2021, 1, 1),
+ start_date=datetime(2024, 1, 1),
catchup=False,
- tags=["example", "auto_ml", "video", "tracking"],
+ tags=["example", "vertex_ai", "auto_ml", "video", "tracking"],
) as dag:
create_bucket = GCSCreateBucketOperator(
task_id="create_bucket",
@@ -106,8 +106,7 @@
project_id=PROJECT_ID,
import_configs=VIDEO_DATA_CONFIG,
)
-
- # [START howto_cloud_create_video_tracking_training_job_operator]
+ # [START how_to_cloud_vertex_ai_create_auto_ml_video_tracking_job_operator]
create_auto_ml_video_training_job = CreateAutoMLVideoTrainingJobOperator(
task_id="auto_ml_video_task",
display_name=VIDEO_DISPLAY_NAME,
@@ -118,7 +117,7 @@
region=REGION,
project_id=PROJECT_ID,
)
- # [END howto_cloud_create_video_tracking_training_job_operator]
+ # [END how_to_cloud_vertex_ai_create_auto_ml_video_tracking_job_operator]
delete_auto_ml_video_training_job = DeleteAutoMLTrainingJobOperator(
task_id="delete_auto_ml_video_training_job",
diff --git a/tests/always/test_project_structure.py b/tests/always/test_project_structure.py
index 85c467160b60e..f12b3ad6a6684 100644
--- a/tests/always/test_project_structure.py
+++ b/tests/always/test_project_structure.py
@@ -354,6 +354,16 @@ class TestGoogleProviderProjectStructure(ExampleCoverageTest, AssetsCoverageTest
"airflow.providers.google.cloud.operators.automl.AutoMLTablesUpdateDatasetOperator",
"airflow.providers.google.cloud.operators.automl.AutoMLDeployModelOperator",
"airflow.providers.google.cloud.operators.automl.AutoMLBatchPredictOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLTrainModelOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLPredictOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLCreateDatasetOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLImportDataOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLGetModelOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLDeleteModelOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLListDatasetOperator",
+ "airflow.providers.google.cloud.operators.automl.AutoMLDeleteDatasetOperator",
+ "airflow.providers.google.cloud.operators.datapipeline.CreateDataPipelineOperator",
+ "airflow.providers.google.cloud.operators.datapipeline.RunDataPipelineOperator",
"airflow.providers.google.cloud.operators.dataproc.DataprocScaleClusterOperator",
"airflow.providers.google.cloud.operators.mlengine.MLEngineManageModelOperator",
"airflow.providers.google.cloud.operators.mlengine.MLEngineManageVersionOperator",