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..d6a2e301d918 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,53 @@ 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 + page = mut.ReadRowsPage(mock_parser, mock_message) + + with mock.patch.dict("sys.modules", {"pandas_gbq": None, "pandas_gbq.arrow": None}): + 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) diff --git a/packages/google-cloud-bigquery/README.rst b/packages/google-cloud-bigquery/README.rst index 705ef5ce712c..56ad7d2d5f28 100644 --- a/packages/google-cloud-bigquery/README.rst +++ b/packages/google-cloud-bigquery/README.rst @@ -21,6 +21,26 @@ processing power of Google's infrastructure. .. _Client Library Documentation: https://cloud.google.com/python/docs/reference/bigquery/latest/summary_overview .. _Product Documentation: https://cloud.google.com/bigquery/docs/reference/v2/ + +Related Libraries +----------------- + +For higher-level data analysis and DataFrame operations, consider using one of the following connector libraries: + +.. list-table:: + :header-rows: 1 + + * - Library + - Description + - Documentation + * - `BigQuery DataFrames (bigframes) `_ + - Flexible pandas-like and scikit-learn-like API powered by BigQuery engine for large-scale data analysis and ML. + - `BigFrames Docs `_ + * - `pandas-gbq `_ + - Convenient integration allowing pandas DataFrames to load from and write to BigQuery. + - `pandas-gbq Docs `_ + + Quick Start ----------- diff --git a/packages/google-cloud-bigquery/noxfile.py b/packages/google-cloud-bigquery/noxfile.py index d4accfea91a4..3d5ff57ee814 100644 --- a/packages/google-cloud-bigquery/noxfile.py +++ b/packages/google-cloud-bigquery/noxfile.py @@ -15,12 +15,12 @@ from __future__ import absolute_import import contextlib +from functools import wraps import os import pathlib import re import shutil import time -from functools import wraps from typing import Generator import nox diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index 9ddec48428b1..04474ac5c387 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2206,7 +2206,7 @@ def test_dbapi_connection_does_not_leak_sockets(self): import gc gc.collect() - for _ in range(30): # Wait up to 3 seconds + for _ in range(60): # Wait up to 6 seconds for background socket cleanup conn_end = current_process.net_connections() conn_count_end = len(conn_end) if conn_count_end <= conn_count_start: diff --git a/packages/pandas-gbq/pandas_gbq/__init__.py b/packages/pandas-gbq/pandas_gbq/__init__.py index a8a54179fd5c..b68c84f18f94 100644 --- a/packages/pandas-gbq/pandas_gbq/__init__.py +++ b/packages/pandas-gbq/pandas_gbq/__init__.py @@ -6,6 +6,7 @@ import sys import warnings +from pandas_gbq import arrow from pandas_gbq import version as pandas_gbq_version from pandas_gbq.contexts import Context, context from pandas_gbq.core.sample import sample @@ -32,4 +33,5 @@ "Context", "context", "sample", + "arrow", ] diff --git a/packages/pandas-gbq/pandas_gbq/arrow.py b/packages/pandas-gbq/pandas_gbq/arrow.py new file mode 100644 index 000000000000..4f38a3e60060 --- /dev/null +++ b/packages/pandas-gbq/pandas_gbq/arrow.py @@ -0,0 +1,44 @@ +"""Arrow integration submodule for pandas-gbq.""" + +from typing import Any, Optional + +try: + import pyarrow as pa +except ImportError: + pa = None # type: ignore[assignment] + + +def from_read_rows_response( + message: Any, + arrow_schema: Optional[Any] = None, +) -> Any: + """Decodes a ReadRowsResponse protobuf message into a pyarrow.RecordBatch.""" + if pa is None: + raise ImportError( + "pyarrow is required to use 'from_read_rows_response'. " + "Please install pyarrow to use this function." + ) + + arrow_record_batch = getattr(message, "arrow_record_batch", None) + if arrow_record_batch is None or not getattr( + arrow_record_batch, "serialized_record_batch", None + ): + empty_schema = arrow_schema or pa.schema([]) + return pa.RecordBatch.from_pylist([], schema=empty_schema) + + serialized_batch = arrow_record_batch.serialized_record_batch + buffer = pa.py_buffer(serialized_batch) + + if arrow_schema is not None: + try: + return pa.ipc.read_record_batch(buffer, arrow_schema) + except Exception: + pass + + try: + reader = pa.ipc.RecordBatchStreamReader(buffer) + return reader.read_next_batch() + except Exception: + msg = pa.ipc.read_message(buffer) + batch_schema = arrow_schema if arrow_schema is not None else pa.schema([]) + return pa.ipc.read_record_batch(msg, batch_schema) diff --git a/packages/pandas-gbq/pandas_gbq/core/pandas.py b/packages/pandas-gbq/pandas_gbq/core/pandas.py index aceaa29b5b91..f656e7a80718 100644 --- a/packages/pandas-gbq/pandas_gbq/core/pandas.py +++ b/packages/pandas-gbq/pandas_gbq/core/pandas.py @@ -3,9 +3,9 @@ # license that can be found in the LICENSE file. import itertools +import typing import pandas -import typing def list_columns_and_indexes(dataframe, index=True): diff --git a/packages/pandas-gbq/pandas_gbq/gbq.py b/packages/pandas-gbq/pandas_gbq/gbq.py index ec6a7e9308aa..b3d2188be206 100644 --- a/packages/pandas-gbq/pandas_gbq/gbq.py +++ b/packages/pandas-gbq/pandas_gbq/gbq.py @@ -63,10 +63,10 @@ def _test_google_api_imports(): try: # google-auth-oauthlib does not have type hints nor stubs that mypy uses for type checking. - # This import is solely to test if the package is installed, so we ignore the "unused import" warning. - # Remove this comment and the ignore pragma upon completing: + # This import is solely to test if the package is installed. + # Remove this comment upon completing: # https://github.com/googleapis/google-cloud-python/issues/17045 - from google_auth_oauthlib.flow import InstalledAppFlow # type: ignore[import-untyped] # noqa: F401 + __import__("google_auth_oauthlib.flow") except ImportError as ex: # pragma: NO COVER raise ImportError("pandas-gbq requires google-auth-oauthlib") from ex diff --git a/packages/pandas-gbq/tests/unit/test_arrow.py b/packages/pandas-gbq/tests/unit/test_arrow.py new file mode 100644 index 000000000000..b18c9c4aee30 --- /dev/null +++ b/packages/pandas-gbq/tests/unit/test_arrow.py @@ -0,0 +1,51 @@ +from unittest import mock + +import pyarrow as pa + +import pandas_gbq.arrow + + +def test_from_read_rows_response_valid_message_returns_record_batch(): + schema = pa.schema([("id", pa.int64()), ("name", pa.string())]) + batch = pa.RecordBatch.from_arrays( + [pa.array([1, 2]), pa.array(["alice", "bob"])], schema=schema + ) + sink = pa.BufferOutputStream() + with pa.ipc.new_stream(sink, schema) as writer: + writer.write_batch(batch) + serialized_bytes = sink.getvalue().to_pybytes() + + mock_message = mock.MagicMock() + mock_message.arrow_record_batch.serialized_record_batch = serialized_bytes + + result_batch = pandas_gbq.arrow.from_read_rows_response( + mock_message, arrow_schema=schema + ) + + assert result_batch.num_rows == 2 + assert result_batch.schema.names == ["id", "name"] + assert result_batch.column(0).to_pylist() == [1, 2] + assert result_batch.column(1).to_pylist() == ["alice", "bob"] + + +def test_from_read_rows_response_empty_message_returns_empty_batch(): + schema = pa.schema([("val", pa.float64())]) + mock_message = mock.MagicMock() + mock_message.arrow_record_batch.serialized_record_batch = b"" + + result_batch = pandas_gbq.arrow.from_read_rows_response( + mock_message, arrow_schema=schema + ) + + assert result_batch.num_rows == 0 + assert result_batch.schema == schema + + +def test_from_read_rows_response_uninstalled_pyarrow_raises_import_error(): + mock_message = mock.MagicMock() + + with mock.patch.object(pandas_gbq.arrow, "pa", None): + import pytest + + with pytest.raises(ImportError, match="pyarrow is required"): + pandas_gbq.arrow.from_read_rows_response(mock_message) diff --git a/packages/pandas-gbq/tests/unit/test_gbq.py b/packages/pandas-gbq/tests/unit/test_gbq.py index 5de78f890e9f..52ebd62cc454 100644 --- a/packages/pandas-gbq/tests/unit/test_gbq.py +++ b/packages/pandas-gbq/tests/unit/test_gbq.py @@ -218,7 +218,9 @@ def test_GbqConnector_get_client_w_new_bq(mock_bigquery_client): connector.get_client() _, kwargs = mock_bigquery_client.call_args - assert kwargs["client_info"].user_agent == "pandas-{}".format(pandas.__version__) + assert kwargs["client_info"].user_agent.startswith( + "pandas-{}".format(pandas.__version__) + ) def test_GbqConnector_process_http_error_transforms_timeout():