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 fbdec9503c7 Support Table objects in BigQueryHook.create_table
table_resource (#73099)
fbdec9503c7 is described below
commit fbdec9503c7f39b32f23508fb9b34c72927755fc
Author: X33834 <[email protected]>
AuthorDate: Tue Oct 6 01:40:07 2026 +0800
Support Table objects in BigQueryHook.create_table table_resource (#73099)
* fix(google): Support Table objects in BigQueryHook.create_table
table_resource
The `table_resource` parameter is typed and documented to accept a dict,
`Table`, `TableReference`, or `TableListItem`, but the implementation only
worked with a plain dict. `Table.from_api_repr()` was called backwards on a
`Table` instance (TypeError) and raw `**`-unpacking was applied to
`TableReference`/`TableListItem` objects (TypeError: not a mapping).
Convert any object input to a plain dict via `.to_api_repr()`, wrapping a
bare `TableReference` under `"tableReference"` since its serialization is
flat, before merging `schema_fields`.
Add unit tests covering dict, Table, TableReference, TableListItem inputs
and schema merging.
Closes: #73091
* Fix table reference attribute in create_table tests
* Add newsfragment for #73099
* bigquery: exercise Table-object table_resource in bigquery_tables system
test
Per review feedback in #73099: add an object-based create_table example
alongside the existing dict-config one, so both supported input types
of BigQueryHook.create_table stay covered by the system test.
* Remove newsfragment
Not needed for provider changes per review feedback in #73099.
* Sort imports in the BigQuery tables example Dag
Generated-by: Claude Opus 5
---------
Co-authored-by: X33834 <[email protected]>
Co-authored-by: Jarek Potiuk <[email protected]>
---
.../providers/google/cloud/hooks/bigquery.py | 13 +++--
.../cloud/bigquery/example_bigquery_tables.py | 18 +++++++
.../tests/unit/google/cloud/hooks/test_bigquery.py | 55 +++++++++++++++++++++-
3 files changed, 82 insertions(+), 4 deletions(-)
diff --git
a/providers/google/src/airflow/providers/google/cloud/hooks/bigquery.py
b/providers/google/src/airflow/providers/google/cloud/hooks/bigquery.py
index e01f752869d..b5f63f5812d 100644
--- a/providers/google/src/airflow/providers/google/cloud/hooks/bigquery.py
+++ b/providers/google/src/airflow/providers/google/cloud/hooks/bigquery.py
@@ -513,11 +513,18 @@ class BigQueryHook(GoogleBaseHook, DbApiHook):
Note that if `retry` is specified, the timeout applies to each
individual attempt.
"""
_table_resource: dict[str, Any] = {}
- if isinstance(table_resource, Table):
- _table_resource = Table.from_api_repr(table_resource) # type:
ignore
+ if isinstance(table_resource, (Table, TableReference, TableListItem)):
+ _table_resource = table_resource.to_api_repr()
+ if isinstance(table_resource, TableReference):
+ # A bare TableReference serializes to a flat dict without the
+ # "tableReference" wrapper expected by the API resource.
+ _table_resource = {"tableReference": _table_resource}
if schema_fields:
_table_resource["schema"] = {"fields": schema_fields}
- table_resource_final = {**table_resource, **_table_resource} # type:
ignore
+ if isinstance(table_resource, dict):
+ table_resource_final = {**table_resource, **_table_resource}
+ else:
+ table_resource_final = _table_resource
table_resource = self._resolve_table_reference(
table_resource=table_resource_final,
project_id=project_id,
diff --git
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
index 4a2814d0590..b82053ed87c 100644
---
a/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
+++
b/providers/google/tests/system/google/cloud/bigquery/example_bigquery_tables.py
@@ -25,6 +25,8 @@ import os
from datetime import datetime
from pathlib import Path
+from google.cloud.bigquery import SchemaField, Table
+
from airflow.models.dag import DAG
from airflow.providers.google.cloud.operators.bigquery import (
BigQueryCreateEmptyDatasetOperator,
@@ -96,6 +98,21 @@ with DAG(
)
# [END howto_operator_bigquery_create_table]
+ # [START howto_operator_bigquery_create_table_from_table_object]
+ create_table_from_table_object = BigQueryCreateTableOperator(
+ task_id="create_table_from_table_object",
+ dataset_id=DATASET_NAME,
+ table_id="test_table_from_object",
+ table_resource=Table(
+ f"{PROJECT_ID}.{DATASET_NAME}.test_table_from_object",
+ schema=[
+ SchemaField("emp_name", "STRING", mode="REQUIRED"),
+ SchemaField("salary", "INTEGER", mode="NULLABLE"),
+ ],
+ ),
+ )
+ # [END howto_operator_bigquery_create_table_from_table_object]
+
# [START howto_operator_bigquery_create_view]
create_view = BigQueryCreateTableOperator(
task_id="create_view",
@@ -258,6 +275,7 @@ with DAG(
# TEST BODY
>> update_dataset
>> create_table
+ >> create_table_from_table_object
>> create_view
>> create_materialized_view
>> update_view
diff --git a/providers/google/tests/unit/google/cloud/hooks/test_bigquery.py
b/providers/google/tests/unit/google/cloud/hooks/test_bigquery.py
index 757c503242c..1e7349aba50 100644
--- a/providers/google/tests/unit/google/cloud/hooks/test_bigquery.py
+++ b/providers/google/tests/unit/google/cloud/hooks/test_bigquery.py
@@ -38,7 +38,7 @@ from google.cloud.bigquery import (
)
from google.cloud.bigquery.dataset import AccessEntry, Dataset, DatasetListItem
from google.cloud.bigquery.routine import Routine
-from google.cloud.bigquery.table import _EmptyRowIterator
+from google.cloud.bigquery.table import TableListItem, _EmptyRowIterator
from google.cloud.exceptions import NotFound
from airflow.exceptions import AirflowProviderDeprecationWarning
@@ -946,6 +946,59 @@ class TestTableOperations(_BigQueryBaseTestClass):
timeout=None,
)
+ @mock.patch("airflow.providers.google.cloud.hooks.bigquery.Client")
+ def test_create_table_with_table_object(self, mock_bq_client):
+ table_resource = Table.from_api_repr({"tableReference":
TABLE_REFERENCE_REPR})
+ self.hook.create_table(
+ project_id=PROJECT_ID,
+ dataset_id=DATASET_ID,
+ table_id=TABLE_ID,
+ table_resource=table_resource,
+ )
+ created_table =
mock_bq_client.return_value.create_table.call_args.kwargs["table"]
+ assert created_table.reference.to_api_repr() == TABLE_REFERENCE_REPR
+
+ @mock.patch("airflow.providers.google.cloud.hooks.bigquery.Client")
+ def test_create_table_with_table_reference_object(self, mock_bq_client):
+ self.hook.create_table(
+ project_id=PROJECT_ID,
+ dataset_id=DATASET_ID,
+ table_id=TABLE_ID,
+ table_resource=TABLE_REFERENCE,
+ )
+ created_table =
mock_bq_client.return_value.create_table.call_args.kwargs["table"]
+ assert created_table.reference.to_api_repr() == TABLE_REFERENCE_REPR
+
+ @mock.patch("airflow.providers.google.cloud.hooks.bigquery.Client")
+ def test_create_table_with_table_list_item(self, mock_bq_client):
+ table_resource = TableListItem({"tableReference":
TABLE_REFERENCE_REPR, "type": "TABLE"})
+ self.hook.create_table(
+ project_id=PROJECT_ID,
+ dataset_id=DATASET_ID,
+ table_id=TABLE_ID,
+ table_resource=table_resource,
+ )
+ created_table =
mock_bq_client.return_value.create_table.call_args.kwargs["table"]
+ assert created_table.reference.to_api_repr() == TABLE_REFERENCE_REPR
+
+ @mock.patch("airflow.providers.google.cloud.hooks.bigquery.Client")
+ def test_create_table_with_table_object_and_schema_fields(self,
mock_bq_client):
+ table_resource = Table.from_api_repr({"tableReference":
TABLE_REFERENCE_REPR})
+ schema_fields = [
+ {"name": "id", "type": "STRING", "mode": "REQUIRED"},
+ {"name": "name", "type": "STRING", "mode": "NULLABLE"},
+ ]
+ self.hook.create_table(
+ project_id=PROJECT_ID,
+ dataset_id=DATASET_ID,
+ table_id=TABLE_ID,
+ table_resource=table_resource,
+ schema_fields=schema_fields,
+ )
+ created_table =
mock_bq_client.return_value.create_table.call_args.kwargs["table"]
+ assert created_table.reference.to_api_repr() == TABLE_REFERENCE_REPR
+ assert created_table.to_api_repr()["schema"]["fields"] == schema_fields
+
@mock.patch("airflow.providers.google.cloud.hooks.bigquery.Client")
def test_get_tables_list(self, mock_client):
table_list = [