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 d2883eead5c Fix Cloud SQL 409 operationInProgress on import/export
operations (#68361)
d2883eead5c is described below
commit d2883eead5c09610901a077aa9c9d127b87adb4c
Author: Radhwène Dhouafli <[email protected]>
AuthorDate: Sun Oct 4 02:43:38 2026 +0200
Fix Cloud SQL 409 operationInProgress on import/export operations (#68361)
* Fix Cloud SQL 409 operationInProgress on import/export operations
Apply the existing operation_in_progress_retry() policy to
CloudSQLHook.import_instance and export_instance, the only two admin
methods that lacked it. Also re-raise operationInProgress HttpError
un-wrapped from import_instance so the retry decorator can see it;
terminal HttpErrors still get the friendly AirflowException message.
* Limit Cloud SQL 409 retry to the import submit call
Address review feedback on #68361: operation_in_progress_retry()
previously wrapped the whole import_instance, including the
operation-status polling. A retryable 429 raised during polling
re-ran the method and re-submitted an import that was already
accepted, importing the same data twice.
The submit now lives in _submit_import, which alone carries the
retry decorator; import_instance waits outside the retry scope, so
a polling failure fails the task instead of re-submitting. Adds a
regression test asserting exactly one submit when polling raises a
retryable error.
* Decode Cloud SQL API error bodies in hook failure messages
* Raise a dedicated error when Cloud SQL import fails
Callers could only tell an import failure apart from any other hook failure
by
matching on the message, since every failure surfaced as the broad
AirflowException. A dedicated subclass lets them catch this case precisely,
and
keeps existing except AirflowException handlers working.
* Keep the Cloud SQL API error when its body is not valid UTF-8
A strict decode raised UnicodeDecodeError while building the
CloudSQLImportError message, replacing the real API error with an unrelated one.
Generated-by: Claude Opus 5
---------
Co-authored-by: Jarek Potiuk <[email protected]>
---
generated/known_airflow_exceptions.txt | 2 +-
.../providers/google/cloud/hooks/cloud_sql.py | 57 +++++++++++++---
.../unit/google/cloud/hooks/test_cloud_sql.py | 78 +++++++++++++++++++++-
3 files changed, 124 insertions(+), 13 deletions(-)
diff --git a/generated/known_airflow_exceptions.txt
b/generated/known_airflow_exceptions.txt
index ca2b46365cb..8fdef403b6b 100644
--- a/generated/known_airflow_exceptions.txt
+++ b/generated/known_airflow_exceptions.txt
@@ -209,7 +209,7 @@
providers/google/src/airflow/providers/google/cloud/hooks/cloud_build.py::3
providers/google/src/airflow/providers/google/cloud/hooks/cloud_composer.py::5
providers/google/src/airflow/providers/google/cloud/hooks/cloud_memorystore.py::5
providers/google/src/airflow/providers/google/cloud/hooks/cloud_run.py::1
-providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py::32
+providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py::31
providers/google/src/airflow/providers/google/cloud/hooks/cloud_storage_transfer_service.py::5
providers/google/src/airflow/providers/google/cloud/hooks/compute.py::6
providers/google/src/airflow/providers/google/cloud/hooks/compute_ssh.py::6
diff --git
a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
index 0f672570f33..4030bd9b63e 100644
--- a/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
+++ b/providers/google/src/airflow/providers/google/cloud/hooks/cloud_sql.py
@@ -65,6 +65,7 @@ from airflow.providers.google.common.hooks.base_google import
(
GoogleBaseAsyncHook,
GoogleBaseHook,
get_field,
+ is_operation_in_progress_exception,
)
from airflow.utils.log.logging_mixin import LoggingMixin
@@ -82,6 +83,10 @@ TIME_TO_SLEEP_IN_SECONDS = 20
CLOUD_SQL_PROXY_VERSION_REGEX = re.compile(r"^v?(\d+\.\d+\.\d+)(-\w*.?\d?)?$")
+class CloudSQLImportError(AirflowException):
+ """Raised when importing data into a Cloud SQL instance fails."""
+
+
class CloudSqlOperationStatus:
"""Helper class with operation statuses."""
@@ -332,10 +337,14 @@ class CloudSQLHook(GoogleBaseHook):
self._wait_for_operation_to_complete(project_id=project_id,
operation_name=operation_name)
@GoogleBaseHook.fallback_to_default_project_id
+ @GoogleBaseHook.operation_in_progress_retry()
def export_instance(self, instance: str, body: dict, project_id: str):
"""
Export data from a Cloud SQL instance to a Cloud Storage bucket as a
SQL dump or CSV file.
+ Cloud SQL runs one admin operation at a time per instance, so this
submit can come back
+ 409 — hence the retry decorator.
+
:param instance: Database instance ID of the Cloud SQL instance. This
does not include the
project ID.
:param body: The request body, as described in
@@ -353,11 +362,45 @@ class CloudSQLHook(GoogleBaseHook):
operation_name = response["name"]
return operation_name
+ @GoogleBaseHook.operation_in_progress_retry()
+ def _submit_import(self, instance: str, body: dict, project_id: str) ->
str:
+ """
+ Submit an import request for a Cloud SQL instance, retrying while an
operation is in progress.
+
+ Cloud SQL runs one admin operation at a time per instance, so this
submit can come back 409.
+ The retry covers the submit only: repeating a rejected submit is safe,
repeating an accepted
+ one is not.
+
+ :param instance: Database instance ID. This does not include the
project ID.
+ :param body: The request body, as described in
+
https://cloud.google.com/sql/docs/mysql/admin-api/v1beta4/instances/import#request-body
+ :param project_id: Project ID of the project that contains the
instance.
+ :return: The name of the accepted import operation.
+ """
+ try:
+ response = (
+ self.get_conn()
+ .instances()
+ .import_(project=project_id, instance=instance, body=body)
+ .execute(num_retries=self.num_retries)
+ )
+ return response["name"]
+ except HttpError as ex:
+ # The decorator retries HttpError, not CloudSQLImportError, so
don't wrap the 409.
+ if is_operation_in_progress_exception(ex):
+ raise
+ raise CloudSQLImportError(
+ f"Importing instance {instance} failed:
{ex.content.decode('utf-8', errors='replace')}"
+ )
+
@GoogleBaseHook.fallback_to_default_project_id
def import_instance(self, instance: str, body: dict, project_id: str) ->
None:
"""
Import data into a Cloud SQL instance from a SQL dump or CSV file in
Cloud Storage.
+ The submit retries on 409 (see ``_submit_import``). The polling below
stays outside that
+ retry: re-submitting an accepted import would load the same data twice.
+
:param instance: Database instance ID. This does not include the
project ID.
:param body: The request body, as described in
@@ -366,17 +409,13 @@ class CloudSQLHook(GoogleBaseHook):
to None or missing, the default project_id from the Google Cloud
connection is used.
:return: None
"""
+ operation_name = self._submit_import(instance=instance, body=body,
project_id=project_id)
try:
- response = (
- self.get_conn()
- .instances()
- .import_(project=project_id, instance=instance, body=body)
- .execute(num_retries=self.num_retries)
- )
- operation_name = response["name"]
self._wait_for_operation_to_complete(project_id=project_id,
operation_name=operation_name)
except HttpError as ex:
- raise AirflowException(f"Importing instance {instance} failed:
{ex.content}")
+ raise CloudSQLImportError(
+ f"Importing instance {instance} failed:
{ex.content.decode('utf-8', errors='replace')}"
+ )
@GoogleBaseHook.fallback_to_default_project_id
def clone_instance(self, instance: str, body: dict, project_id: str) ->
None:
@@ -403,7 +442,7 @@ class CloudSQLHook(GoogleBaseHook):
operation_name = response["name"]
self._wait_for_operation_to_complete(project_id=project_id,
operation_name=operation_name)
except HttpError as ex:
- raise AirflowException(f"Cloning of instance {instance} failed:
{ex.content}")
+ raise AirflowException(f"Cloning of instance {instance} failed:
{ex.content.decode('utf-8')}")
@GoogleBaseHook.fallback_to_default_project_id
def create_ssl_certificate(self, instance: str, body: dict, project_id:
str):
diff --git a/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
b/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
index 955add00379..07e2f9be6e8 100644
--- a/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
+++ b/providers/google/tests/unit/google/cloud/hooks/test_cloud_sql.py
@@ -45,6 +45,7 @@ from airflow.providers.google.cloud.hooks.cloud_sql import (
CloudSQLAsyncHook,
CloudSQLDatabaseHook,
CloudSQLHook,
+ CloudSQLImportError,
CloudSqlProxyRunner,
)
@@ -96,7 +97,7 @@ class TestGcpSqlHookDefaultProjectId:
self.cloudsql_hook.get_conn = mock.Mock(
side_effect=HttpError(resp=httplib2.Response({"status": 400}),
content=b"Error content")
)
- with pytest.raises(AirflowException) as ctx:
+ with pytest.raises(CloudSQLImportError) as ctx:
self.cloudsql_hook.import_instance(instance="instance", body={})
err = ctx.value
assert "Importing instance " in str(err)
@@ -165,10 +166,81 @@ class TestGcpSqlHookDefaultProjectId:
),
{"name": "operation_id"},
]
- with pytest.raises(HttpError):
- self.cloudsql_hook.export_instance(project_id="example-project",
instance="instance", body={})
+ # First submit returns 429 (one of the two operation-in-progress codes
recognised by
+ # ``is_operation_in_progress_exception``);
``operation_in_progress_retry`` retries and the
+ # second submit succeeds, returning the operation name. The import
test below covers 409.
+ result = self.cloudsql_hook.export_instance(
+ project_id="example-project", instance="instance", body={}
+ )
+ assert result == "operation_id"
+ assert export_method.call_count == 2
+ assert execute_method.call_count == 2
wait_for_operation_to_complete.assert_not_called()
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+ def test_instance_import_with_in_progress_retry(self,
wait_for_operation_to_complete, get_conn):
+ import_method = get_conn.return_value.instances.return_value.import_
+ execute_method = import_method.return_value.execute
+ execute_method.side_effect = [
+ HttpError(
+ resp=httplib2.Response({"status": 409}),
+ content=b"operationInProgress",
+ ),
+ {"name": "operation_id"},
+ ]
+ wait_for_operation_to_complete.return_value = None
+ # First submit returns 409 ``operationInProgress``. ``_submit_import``
re-raises it past its
+ # friendly-message wrapper (instead of converting it to
CloudSQLImportError), so
+ # ``operation_in_progress_retry`` sees the raw HttpError, retries, and
the second submit
+ # succeeds; the resulting operation is awaited exactly once.
+ self.cloudsql_hook.import_instance(project_id="example-project",
instance="instance", body={})
+ assert import_method.call_count == 2
+ assert execute_method.call_count == 2
+ wait_for_operation_to_complete.assert_called_once_with(
+ project_id="example-project", operation_name="operation_id"
+ )
+
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+ def test_instance_import_does_not_resubmit_when_polling_fails(
+ self, wait_for_operation_to_complete, get_conn
+ ):
+ import_method = get_conn.return_value.instances.return_value.import_
+ execute_method = import_method.return_value.execute
+ execute_method.return_value = {"name": "operation_id"}
+ wait_for_operation_to_complete.side_effect = HttpError(
+ resp=httplib2.Response({"status": 429}),
+ content=b"rate limited",
+ )
+ # The submit succeeds, then the operation-status polling raises a
retryable 429. The retry
+ # scope must not include the polling: re-running ``import_instance``
would re-submit an
+ # import that was already accepted and import the same data twice. The
task must fail
+ # instead, with exactly one submit on record.
+ with pytest.raises(CloudSQLImportError, match="Importing instance
instance failed"):
+ self.cloudsql_hook.import_instance(project_id="example-project",
instance="instance", body={})
+ import_method.assert_called_once_with(body={}, instance="instance",
project="example-project")
+ execute_method.assert_called_once_with(num_retries=5)
+ wait_for_operation_to_complete.assert_called_once_with(
+ project_id="example-project", operation_name="operation_id"
+ )
+
+ @pytest.mark.parametrize("failing_step", ["submit", "polling"])
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_conn")
+
@mock.patch("airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook._wait_for_operation_to_complete")
+ def test_instance_import_error_with_non_utf8_body(
+ self, wait_for_operation_to_complete, get_conn, failing_step
+ ):
+ error = HttpError(resp=httplib2.Response({"status": 400}),
content=b"bad \xff body")
+ execute_method =
get_conn.return_value.instances.return_value.import_.return_value.execute
+ if failing_step == "submit":
+ execute_method.side_effect = error
+ else:
+ execute_method.return_value = {"name": "operation_id"}
+ wait_for_operation_to_complete.side_effect = error
+ with pytest.raises(CloudSQLImportError, match="Importing instance
instance failed: bad � body"):
+ self.cloudsql_hook.import_instance(project_id="example-project",
instance="instance", body={})
+
@mock.patch(
"airflow.providers.google.cloud.hooks.cloud_sql.CloudSQLHook.get_credentials_and_project_id",
return_value=(mock.MagicMock(), "example-project"),