Skip to content

Commit 67cc6c2

Browse files
committed
feat(bigquery): support Arrow destination tables
BigQuery Storage Read API streams/_default only supports anonymous query result tables and rejects explicit destination tables with INVALID_ARGUMENT. When query_and_wait falls back to jobs.insert with an explicit destination table, open a single-stream Arrow read session via create_read_session (passing compression_codec in read_options) and stream batches over read_rows while preserving 0-row Arrow schemas.
1 parent b0639a7 commit 67cc6c2

4 files changed

Lines changed: 150 additions & 4 deletions

File tree

‎packages/google-cloud-bigquery/google/cloud/bigquery/_job_helpers.py‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -517,7 +517,7 @@ def query_and_wait(
517517
# Some API parameters aren't supported by the jobs.query API. In these
518518
# cases, fallback to a jobs.insert call.
519519
if not _supported_by_jobs_query(request_body):
520-
return _wait_or_cancel(
520+
rows = _wait_or_cancel(
521521
query_jobs_insert(
522522
client=client,
523523
query=query,
@@ -540,6 +540,9 @@ def query_and_wait(
540540
compression_codec=compression_codec,
541541
callback=callback,
542542
)
543+
if "destinationTable" in request_body:
544+
rows._use_read_session = True
545+
return rows
543546

544547
path = _to_query_path(project)
545548

‎packages/google-cloud-bigquery/google/cloud/bigquery/table.py‎

Lines changed: 39 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2374,10 +2374,46 @@ def _stream_arrow_via_bqstorage(
23742374
remaining = (
23752375
self.max_results - offset if self.max_results is not None else None
23762376
)
2377-
location = self._location or (self.client.location if self.client else None)
2378-
stream_name = f"projects/{project}/locations/{location}/jobs/{self._job_id}/streams/_default"
2377+
use_read_session = (
2378+
getattr(self, "_use_read_session", False) and self._table is not None
2379+
)
2380+
if use_read_session:
2381+
table_ref = (
2382+
f"projects/{self._table.project}/datasets/"
2383+
f"{self._table.dataset_id}/tables/{self._table.table_id}"
2384+
)
2385+
read_session: Dict[str, Any] = {
2386+
"table": table_ref,
2387+
"data_format": "ARROW",
2388+
}
2389+
if self._compression_codec is not None:
2390+
read_session["read_options"] = {
2391+
"arrow_serialization_options": {
2392+
"buffer_compression": self._compression_codec
2393+
}
2394+
}
2395+
session = bqstorage_client.create_read_session(
2396+
parent=f"projects/{project}",
2397+
read_session=read_session,
2398+
max_stream_count=1,
2399+
)
2400+
if (
2401+
session.arrow_schema
2402+
and session.arrow_schema.serialized_schema
2403+
and pa_schema is None
2404+
):
2405+
pa_schema = pyarrow.ipc.read_schema(
2406+
pyarrow.py_buffer(session.arrow_schema.serialized_schema)
2407+
)
2408+
self._arrow_schema = pa_schema
2409+
if not session.streams:
2410+
return
2411+
stream_name = session.streams[0].name
2412+
else:
2413+
location = self._location or (self.client.location if self.client else None)
2414+
stream_name = f"projects/{project}/locations/{location}/jobs/{self._job_id}/streams/_default"
23792415
read_kwargs: Dict[str, Any] = {"offset": offset, "timeout": timeout}
2380-
if self._compression_codec is not None:
2416+
if self._compression_codec is not None and not use_read_session:
23812417
try:
23822418
reader = bqstorage_client.read_rows(
23832419
stream_name,

‎packages/google-cloud-bigquery/tests/system/test_arrow.py‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -474,3 +474,44 @@ def test_query_arrow_cached_multi_page(bigquery_client):
474474
assert len(first_table) == 5000
475475
assert len(cached_table) == 5000
476476
assert cached_table.equals(first_table)
477+
478+
479+
def test_query_arrow_explicit_destination_table(
480+
dataset_client, dataset_id, test_table_name
481+
):
482+
"""System test verifying explicit destination tables on QueryJobConfig read via CreateReadSession."""
483+
dest_table = dataset_client.dataset(dataset_id).table(test_table_name)
484+
job_config = bigquery.QueryJobConfig(
485+
destination=dest_table,
486+
write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE,
487+
)
488+
489+
# 1. Full + max_results read with LZ4_FRAME compression on explicit destination table
490+
results = dataset_client.query_and_wait(
491+
"SELECT num, CONCAT('row_', CAST(num AS STRING)) AS label "
492+
"FROM UNNEST(GENERATE_ARRAY(1, 2000)) AS num ORDER BY num",
493+
job_config=job_config,
494+
query_results_format=enums.QueryResultsFormat.ARROW,
495+
compression_codec=enums.QueryResultsCompressionCodec.LZ4_FRAME,
496+
max_results=1200,
497+
)
498+
table = pyarrow.Table.from_batches(results.to_arrow_iterable())
499+
assert len(table) == 1200
500+
assert table.column_names == ["num", "label"]
501+
assert table.column("num").to_pylist() == list(range(1, 1201))
502+
503+
# 2. Zero-row explicit destination table preserves pyarrow.Schema
504+
empty_dest_table = dataset_client.dataset(dataset_id).table(
505+
f"{test_table_name}_empty"
506+
)
507+
empty_results = dataset_client.query_and_wait(
508+
"SELECT 1 AS num, 'abc' AS label FROM UNNEST(GENERATE_ARRAY(1, 5)) AS x WHERE x > 100",
509+
job_config=bigquery.QueryJobConfig(destination=empty_dest_table),
510+
query_results_format=enums.QueryResultsFormat.ARROW,
511+
)
512+
assert empty_results.total_rows == 0
513+
empty_table = empty_results.to_arrow()
514+
assert len(empty_table) == 0
515+
assert empty_table.column_names == ["num", "label"]
516+
assert empty_table.schema.field("num").type == pyarrow.int64()
517+
assert empty_table.schema.field("label").type == pyarrow.string()

‎packages/google-cloud-bigquery/tests/unit/test_query_results_format_arrow.py‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -914,6 +914,72 @@ def test_download_arrow_from_job_id_cached_query_with_page_token_reads_stream(se
914914
timeout=5.0,
915915
)
916916

917+
def test_download_arrow_from_job_id_explicit_destination_uses_create_read_session(
918+
self,
919+
):
920+
from google.cloud.bigquery.table import TableReference
921+
922+
mock_client = mock.MagicMock()
923+
mock_bqstorage = mock.MagicMock()
924+
mock_client._ensure_bqstorage_client.return_value = mock_bqstorage
925+
926+
mock_stream = mock.MagicMock()
927+
mock_stream.name = "projects/test-proj/locations/US/sessions/s1/streams/0"
928+
mock_session = mock.MagicMock()
929+
mock_session.arrow_schema.serialized_schema = b"session_schema_bytes"
930+
mock_session.streams = [mock_stream]
931+
mock_bqstorage.create_read_session.return_value = mock_session
932+
933+
mock_resp = mock.MagicMock()
934+
mock_resp.arrow_schema = None
935+
mock_resp.arrow_record_batch.serialized_record_batch = b"batch_bytes"
936+
mock_bqstorage.read_rows.return_value = [mock_resp]
937+
938+
dest_table = TableReference.from_string("test-proj.dest_ds.dest_tbl")
939+
iterator = RowIterator(
940+
client=mock_client,
941+
api_request=mock.MagicMock(),
942+
path=None,
943+
schema=(),
944+
table=dest_table,
945+
project="test-proj",
946+
location="US",
947+
job_id="test-job-dest",
948+
query_results_format="ARROW",
949+
compression_codec="LZ4_FRAME",
950+
)
951+
iterator._use_read_session = True
952+
953+
mock_batch = mock.MagicMock()
954+
mock_batch.num_rows = 3
955+
956+
with mock.patch("google.cloud.bigquery.table.pyarrow") as mock_pyarrow:
957+
mock_pyarrow.py_buffer = lambda x: x
958+
mock_pyarrow.ipc.read_schema.return_value = "session_schema"
959+
mock_pyarrow.ipc.read_record_batch.return_value = mock_batch
960+
961+
batches = list(iterator._download_arrow_from_job_id(timeout=5.0))
962+
self.assertEqual(batches, [mock_batch])
963+
self.assertEqual(iterator._arrow_schema, "session_schema")
964+
mock_bqstorage.create_read_session.assert_called_once_with(
965+
parent="projects/test-proj",
966+
read_session={
967+
"table": "projects/test-proj/datasets/dest_ds/tables/dest_tbl",
968+
"data_format": "ARROW",
969+
"read_options": {
970+
"arrow_serialization_options": {
971+
"buffer_compression": "LZ4_FRAME"
972+
}
973+
},
974+
},
975+
max_stream_count=1,
976+
)
977+
mock_bqstorage.read_rows.assert_called_once_with(
978+
"projects/test-proj/locations/US/sessions/s1/streams/0",
979+
offset=0,
980+
timeout=5.0,
981+
)
982+
917983

918984
if __name__ == "__main__":
919985
unittest.main()

0 commit comments

Comments
 (0)