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 e6341f3e748 Delay BigQueryToMsSqlOperator source table parsing (#70493)
e6341f3e748 is described below
commit e6341f3e7483cc4451d8443434a9980f1f1f9341
Author: Vincent Hsiao <[email protected]>
AuthorDate: Sat Aug 1 09:29:29 2026 +0800
Delay BigQueryToMsSqlOperator source table parsing (#70493)
* Delay BigQueryToMsSqlOperator source table parsing
Template fields are rendered after operator construction, so parsing the
BigQuery source table in the constructor rejects valid Jinja expressions before
tasks can run.
* Avoid placeholder BigQuery table parts in MSSQL transfer
The MSSQL transfer derives source table parts from a templated
project.dataset.table value at execution time, so initializing the shared base
class with synthetic table parts can expose confusing internal placeholders
before rendering completes.
---
.../google/cloud/transfers/bigquery_to_mssql.py | 22 +++++----
.../google/cloud/transfers/bigquery_to_sql.py | 11 +++--
.../cloud/transfers/test_bigquery_to_mssql.py | 53 ++++++++++++++++++++++
.../ci/prek/validate_operators_init_exemptions.txt | 1 -
4 files changed, 73 insertions(+), 14 deletions(-)
diff --git
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_mssql.py
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_mssql.py
index d2ebe9737ec..93b12146c3f 100644
---
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_mssql.py
+++
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_mssql.py
@@ -81,15 +81,8 @@ class BigQueryToMsSqlOperator(BigQueryToSqlBaseOperator):
target_table_name = mssql_table
- try:
- _, dataset_id, table_id = source_project_dataset_table.split(".")
- except ValueError:
- raise ValueError(
- f"Could not parse {source_project_dataset_table} as
<project>.<dataset>.<table>"
- ) from None
super().__init__(
target_table_name=target_table_name,
- dataset_table=f"{dataset_id}.{table_id}",
**kwargs,
)
self.mssql_conn_id = mssql_conn_id
@@ -102,8 +95,21 @@ class BigQueryToMsSqlOperator(BigQueryToSqlBaseOperator):
def get_sql_hook(self) -> MsSqlHook:
return self.mssql_hook
+ def execute(self, context: Context) -> None:
+ _, self.dataset_id, self.table_id =
self._get_source_project_dataset_table_parts()
+ super().execute(context)
+
+ def _get_source_project_dataset_table_parts(self) -> tuple[str, str, str]:
+ try:
+ project_id, dataset_id, table_id =
self.source_project_dataset_table.split(".")
+ except ValueError:
+ raise ValueError(
+ f"Could not parse {self.source_project_dataset_table} as
<project>.<dataset>.<table>"
+ ) from None
+ return project_id, dataset_id, table_id
+
def persist_links(self, context: Context) -> None:
- project_id, dataset_id, table_id =
self.source_project_dataset_table.split(".")
+ project_id, dataset_id, table_id =
self._get_source_project_dataset_table_parts()
BigQueryTableLink.persist(
context=context,
dataset_id=dataset_id,
diff --git
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
index 81218f0aa97..5c188fe8d04 100644
---
a/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
+++
b/providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_sql.py
@@ -79,8 +79,8 @@ class BigQueryToSqlBaseOperator(BaseOperator):
def __init__(
self,
*,
- dataset_table: str,
target_table_name: str | None,
+ dataset_table: str | None = None,
selected_fields: list[str] | str | None = None,
gcp_conn_id: str = "google_cloud_default",
database: str | None = None,
@@ -101,12 +101,13 @@ class BigQueryToSqlBaseOperator(BaseOperator):
self.batch_size = batch_size
self.location = location
self.impersonation_chain = impersonation_chain
+ if dataset_table is not None:
+ try:
+ dataset_id, table_id = dataset_table.split(".")
+ except ValueError:
+ raise ValueError(f"Could not parse {dataset_table} as
<dataset>.<table>") from None
self.dataset_id = dataset_id
self.table_id = table_id
- try:
- self.dataset_id, self.table_id = dataset_table.split(".")
- except ValueError:
- raise ValueError(f"Could not parse {dataset_table} as
<dataset>.<table>") from None
@abc.abstractmethod
def get_sql_hook(self) -> DbApiHook:
diff --git
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
index a9c4668b50e..0295f7060f2 100644
---
a/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
+++
b/providers/google/tests/unit/google/cloud/transfers/test_bigquery_to_mssql.py
@@ -20,6 +20,8 @@ from __future__ import annotations
from unittest import mock
from unittest.mock import MagicMock
+import pytest
+
from airflow.providers.google.cloud.transfers.bigquery_to_mssql import
BigQueryToMsSqlOperator
TASK_ID = "test-bq-create-table-operator"
@@ -68,6 +70,57 @@ class TestBigQueryToMsSqlOperator:
start_index=0,
)
+
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_mssql.BigQueryTableLink")
+
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_sql.BigQueryHook")
+ def test_execute_uses_rendered_source_project_dataset_table(self,
mock_hook, mock_link):
+ destination_table = "table"
+ source_project_dataset_table =
f"{TEST_PROJECT}.{TEST_DATASET}.{TEST_TABLE_ID}"
+ operator = BigQueryToMsSqlOperator(
+ task_id=TASK_ID,
+ source_project_dataset_table="{{
var.value.source_project_dataset_table }}",
+ target_table_name=destination_table,
+ replace=False,
+ )
+
+ assert operator.dataset_id is None
+ assert operator.table_id is None
+
+ operator.render_template_fields(
+ {"var": {"value": {"source_project_dataset_table":
source_project_dataset_table}}}
+ )
+ operator.execute(None)
+
+ mock_hook.return_value.list_rows.assert_called_once_with(
+ dataset_id=TEST_DATASET,
+ table_id=TEST_TABLE_ID,
+ max_results=1000,
+ selected_fields=None,
+ start_index=0,
+ )
+ mock_link.persist.assert_called_once_with(
+ context=None,
+ dataset_id=TEST_DATASET,
+ project_id=TEST_PROJECT,
+ table_id=TEST_TABLE_ID,
+ )
+
+ def
test_execute_raises_for_invalid_rendered_source_project_dataset_table(self):
+ operator = BigQueryToMsSqlOperator(
+ task_id=TASK_ID,
+ source_project_dataset_table="{{
var.value.source_project_dataset_table }}",
+ target_table_name="table",
+ replace=False,
+ )
+
+ operator.render_template_fields(
+ {"var": {"value": {"source_project_dataset_table":
"project.dataset"}}}
+ )
+
+ with pytest.raises(
+ ValueError, match="Could not parse project.dataset as
<project>.<dataset>.<table>"
+ ):
+ operator.execute(None)
+
@mock.patch("airflow.providers.google.cloud.transfers.bigquery_to_mssql.MsSqlHook")
def test_get_sql_hook(self, mock_hook):
hook_expected = mock_hook.return_value
diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt
b/scripts/ci/prek/validate_operators_init_exemptions.txt
index a41b768d9be..e6869670a5b 100644
--- a/scripts/ci/prek/validate_operators_init_exemptions.txt
+++ b/scripts/ci/prek/validate_operators_init_exemptions.txt
@@ -19,7 +19,6 @@
providers/google/src/airflow/providers/google/cloud/operators/gcs.py::GCSFileTra
providers/google/src/airflow/providers/google/cloud/sensors/bigquery_dts.py::BigQueryDataTransferServiceTransferRunSensor
providers/google/src/airflow/providers/google/cloud/sensors/cloud_composer.py::CloudComposerExternalTaskSensor
providers/google/src/airflow/providers/google/cloud/transfers/azure_fileshare_to_gcs.py::AzureFileShareToGCSOperator
-providers/google/src/airflow/providers/google/cloud/transfers/bigquery_to_mssql.py::BigQueryToMsSqlOperator
providers/google/src/airflow/providers/google/cloud/transfers/gcs_to_bigquery.py::GCSToBigQueryOperator
providers/google/src/airflow/providers/google/marketing_platform/operators/campaign_manager.py::GoogleCampaignManagerDeleteReportOperator
providers/microsoft/azure/src/airflow/providers/microsoft/azure/transfers/gcs_to_wasb.py::GCSToAzureBlobStorageOperator