From a63af8f406d0dc5a0b89b25f2949abb3f39ddd37 Mon Sep 17 00:00:00 2001 From: Shuowei Li Date: Thu, 30 Jul 2026 21:24:19 +0000 Subject: [PATCH] feat(storage): delegate ReadRowsPage.to_arrow to pandas-gbq --- .../cloud/bigquery_storage_v1/reader.py | 20 ++++++++ .../tests/unit/test_reader_v1_arrow.py | 49 +++++++++++++++++++ 2 files changed, 69 insertions(+) diff --git a/packages/google-cloud-bigquery-storage/google/cloud/bigquery_storage_v1/reader.py b/packages/google-cloud-bigquery-storage/google/cloud/bigquery_storage_v1/reader.py index c4caa536da1d..9ef2c78d702a 100644 --- a/packages/google-cloud-bigquery-storage/google/cloud/bigquery_storage_v1/reader.py +++ b/packages/google-cloud-bigquery-storage/google/cloud/bigquery_storage_v1/reader.py @@ -18,6 +18,7 @@ import io import json import time +import warnings try: import fastavro @@ -572,6 +573,25 @@ def to_arrow(self): pyarrow.RecordBatch: Rows from the message, as an Arrow record batch. """ + warnings.warn( + "Retrieving Arrow record batches directly via google-cloud-bigquery-storage is deprecated. " + "Please use 'pandas_gbq.arrow.from_read_rows_response' or install 'pandas-gbq'.", + PendingDeprecationWarning, + stacklevel=2, + ) + try: + import pandas_gbq.arrow # type: ignore[import-not-found] + + if hasattr(pandas_gbq.arrow, "from_read_rows_response"): + if hasattr(self._stream_parser, "_parse_arrow_schema"): + self._stream_parser._parse_arrow_schema() + arrow_schema = getattr(self._stream_parser, "_schema", None) + return pandas_gbq.arrow.from_read_rows_response( + self._message, arrow_schema=arrow_schema + ) + except ImportError: + pass + return self._stream_parser.to_arrow(self._message) def to_dataframe(self, dtypes=None): diff --git a/packages/google-cloud-bigquery-storage/tests/unit/test_reader_v1_arrow.py b/packages/google-cloud-bigquery-storage/tests/unit/test_reader_v1_arrow.py index ca9f348351db..7541f0080e31 100644 --- a/packages/google-cloud-bigquery-storage/tests/unit/test_reader_v1_arrow.py +++ b/packages/google-cloud-bigquery-storage/tests/unit/test_reader_v1_arrow.py @@ -428,3 +428,52 @@ def test_to_dataframe_mid_stream_failure(mut, class_under_test, mock_gapic_clien with pytest.raises(RuntimeError, match="Stream format changed mid-stream"): it.to_dataframe() + + +def test_to_arrow_delegates_to_pandas_gbq_when_installed(mut): + mock_parser = mock.Mock() + mock_message = mock.Mock() + page = mut.ReadRowsPage(mock_parser, mock_message) + expected_batch = pyarrow.RecordBatch.from_arrays( + [pyarrow.array([100])], names=["id"] + ) + + mock_arrow_module = mock.Mock() + mock_arrow_module.from_read_rows_response.return_value = expected_batch + + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.arrow = mock_arrow_module + + with mock.patch.dict( + "sys.modules", + {"pandas_gbq": mock_pandas_gbq, "pandas_gbq.arrow": mock_arrow_module}, + ): + with pytest.warns( + PendingDeprecationWarning, + match="google-cloud-bigquery-storage is deprecated", + ): + actual_batch = page.to_arrow() + + assert actual_batch == expected_batch + mock_arrow_module.from_read_rows_response.assert_called_once_with( + mock_message, arrow_schema=mock_parser._schema + ) + + +def test_to_arrow_falls_back_when_pandas_gbq_uninstalled(mut): + mock_parser = mock.Mock() + mock_message = mock.Mock() + expected_batch = pyarrow.RecordBatch.from_arrays( + [pyarrow.array([200])], names=["id"] + ) + mock_parser.to_arrow.return_value = expected_batch + with mock.patch.dict("sys.modules", {"pandas_gbq": None, "pandas_gbq.arrow": None}): + page = mut.ReadRowsPage(mock_parser, mock_message) + with pytest.warns( + PendingDeprecationWarning, + match="google-cloud-bigquery-storage is deprecated", + ): + actual_batch = page.to_arrow() + + assert actual_batch == expected_batch + mock_parser.to_arrow.assert_called_once_with(mock_message)