This is an automated email from the ASF dual-hosted git repository.
potiuk 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 08fae625769 Add datafusion handler for google cleanup script (#74234)
08fae625769 is described below
commit 08fae625769ed5d4874fd695bdb415b951c6f472
Author: Ulada Zakharava <[email protected]>
AuthorDate: Mon Oct 5 18:00:38 2026 +0200
Add datafusion handler for google cleanup script (#74234)
---
.../system/google/resources_cleanup/README.md | 8 +++
.../handlers/__init__.py | 2 +
.../handlers/datafusion.py | 52 +++++++++++++++++
.../google/resources_cleanup/test_cmd_delete.py | 1 +
.../resources_cleanup/test_datafusion_handler.py | 67 ++++++++++++++++++++++
5 files changed, 130 insertions(+)
diff --git a/providers/google/tests/system/google/resources_cleanup/README.md
b/providers/google/tests/system/google/resources_cleanup/README.md
index dd44cafb383..e0e45fff0e0 100644
--- a/providers/google/tests/system/google/resources_cleanup/README.md
+++ b/providers/google/tests/system/google/resources_cleanup/README.md
@@ -111,6 +111,14 @@ Example configuration:
- To clean up only resources that are old enough:
`python -m airflow_google_provider_resource_cleanup delete --project-id
<PROJECT_ID> --asset-type dataproc --min-age-days 3`
+- To clean up Cloud Data Fusion DNS peerings and instances:
+ `python -m airflow_google_provider_resource_cleanup delete --project-id
<PROJECT_ID> --asset-type datafusion`
+
+ Cloud Data Fusion cleanup processes DNS peerings before their parent
+ instances. Instance deletion does not use the API's `force` option, so an
+ instance with remaining nested resources is left in place instead of
+ bypassing resource protection or minimum-age filtering.
+
- To skip a service group during deletion, for example Composer:
`python -m airflow_google_provider_resource_cleanup delete --project-id
<PROJECT_ID> --skip-asset-type composer`
diff --git
a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/__init__.py
b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/__init__.py
index d3ed1676b9e..f8225dd4a8b 100644
---
a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/__init__.py
+++
b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/__init__.py
@@ -30,6 +30,7 @@ def get_delete_handlers() -> dict[str,
type[BaseDeleteHandler]]:
from airflow_google_provider_resource_cleanup.handlers.compute import
ComputeDeleteHandler
from airflow_google_provider_resource_cleanup.handlers.dataflow import
DataflowDeleteHandler
from airflow_google_provider_resource_cleanup.handlers.dataform import
DataformDeleteHandler
+ from airflow_google_provider_resource_cleanup.handlers.datafusion import
DataFusionDeleteHandler
from airflow_google_provider_resource_cleanup.handlers.dataplex import
DataplexDeleteHandler
from airflow_google_provider_resource_cleanup.handlers.dataproc import
DataprocDeleteHandler
from airflow_google_provider_resource_cleanup.handlers.dataproc_metastore
import (
@@ -53,6 +54,7 @@ def get_delete_handlers() -> dict[str,
type[BaseDeleteHandler]]:
"dataflow": DataflowDeleteHandler,
"sqladmin": CloudSQLDeleteHandler,
"dataform": DataformDeleteHandler,
+ "datafusion": DataFusionDeleteHandler,
"dataplex": DataplexDeleteHandler,
"dataproc": DataprocDeleteHandler,
"dataproc_metastore": DataprocMetastoreDeleteHandler,
diff --git
a/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/datafusion.py
b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/datafusion.py
new file mode 100644
index 00000000000..9b867314f82
--- /dev/null
+++
b/providers/google/tests/system/google/resources_cleanup/airflow_google_provider_resource_cleanup/handlers/datafusion.py
@@ -0,0 +1,52 @@
+#
+# 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.
+from __future__ import annotations
+
+from airflow_google_provider_resource_cleanup.handlers._base import
BaseDeleteHandler
+from airflow_google_provider_resource_cleanup.helpers import curl,
get_resource_path, run_command_async
+
+API_BASE = "https://datafusion.googleapis.com/v1/"
+
+
+async def _delete_dns_peering(resource: dict, log_prefix: str):
+ url = f"{API_BASE}{get_resource_path(resource)}"
+ await curl(url, log_prefix=log_prefix)
+
+
+async def _delete_instance(resource: dict, log_prefix: str):
+ path_parts = get_resource_path(resource).split("/")
+ project_id = path_parts[1]
+ location = path_parts[3]
+ instance_id = path_parts[-1]
+ cmd = (
+ f"gcloud beta data-fusion instances delete {instance_id} "
+ f"--location={location} --project={project_id} --quiet"
+ )
+ await run_command_async(cmd, log_prefix)
+
+
+class DataFusionDeleteHandler(BaseDeleteHandler):
+ DELETERS = {
+ "datafusion.googleapis.com/DnsPeering": _delete_dns_peering,
+ "datafusion.googleapis.com/Instance": _delete_instance,
+ }
+
+ DELETION_ORDER = [
+ "datafusion.googleapis.com/DnsPeering",
+ "datafusion.googleapis.com/Instance",
+ ]
diff --git
a/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py
b/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py
index bc34041d9d4..13c951f64bb 100644
--- a/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py
+++ b/providers/google/tests/unit/google/resources_cleanup/test_cmd_delete.py
@@ -176,6 +176,7 @@ def
test_get_delete_handlers_registers_implemented_handlers():
for asset_type in [
"compute",
"dataflow",
+ "datafusion",
"dataproc_metastore",
"gke",
"kafka",
diff --git
a/providers/google/tests/unit/google/resources_cleanup/test_datafusion_handler.py
b/providers/google/tests/unit/google/resources_cleanup/test_datafusion_handler.py
new file mode 100644
index 00000000000..9d47583d306
--- /dev/null
+++
b/providers/google/tests/unit/google/resources_cleanup/test_datafusion_handler.py
@@ -0,0 +1,67 @@
+#
+# 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.
+from __future__ import annotations
+
+from unittest.mock import patch
+
+import pytest
+from airflow_google_provider_resource_cleanup.handlers import datafusion
+
+
[email protected]
[email protected](datafusion, "curl", autospec=True)
+async def test_delete_dns_peering(mock_curl):
+ resource_name = (
+
"//datafusion.googleapis.com/projects/test-project/locations/us-central1/instances/"
+ "test-instance/dnsPeerings/test-peering"
+ )
+
+ await
datafusion.DataFusionDeleteHandler.DELETERS["datafusion.googleapis.com/DnsPeering"](
+ {"name": resource_name}, "[1/1] "
+ )
+
+ mock_curl.assert_awaited_once_with(
+
"https://datafusion.googleapis.com/v1/projects/test-project/locations/us-central1/instances/"
+ "test-instance/dnsPeerings/test-peering",
+ log_prefix="[1/1] ",
+ )
+
+
[email protected]
[email protected](datafusion, "run_command_async", autospec=True)
+async def test_delete_instance(mock_run_command_async):
+ resource_name = (
+
"//datafusion.googleapis.com/projects/test-project/locations/us-central1/instances/test-instance"
+ )
+
+ await
datafusion.DataFusionDeleteHandler.DELETERS["datafusion.googleapis.com/Instance"](
+ {"name": resource_name}, "[1/1] "
+ )
+
+ mock_run_command_async.assert_awaited_once_with(
+ "gcloud beta data-fusion instances delete test-instance "
+ "--location=us-central1 --project=test-project --quiet",
+ "[1/1] ",
+ )
+
+
+def test_deletion_order_removes_dns_peerings_before_instances():
+ assert datafusion.DataFusionDeleteHandler.DELETION_ORDER == [
+ "datafusion.googleapis.com/DnsPeering",
+ "datafusion.googleapis.com/Instance",
+ ]