diff --git a/packages/google-cloud-bigquery/google/cloud/bigquery/_versions_helpers.py b/packages/google-cloud-bigquery/google/cloud/bigquery/_versions_helpers.py index d856c19852e7..ecc4e22e9a0d 100644 --- a/packages/google-cloud-bigquery/google/cloud/bigquery/_versions_helpers.py +++ b/packages/google-cloud-bigquery/google/cloud/bigquery/_versions_helpers.py @@ -16,10 +16,8 @@ from typing import Any import packaging.version - from google.cloud.bigquery import exceptions - _MIN_PYARROW_VERSION = packaging.version.Version("3.0.0") _MIN_BQ_STORAGE_VERSION = packaging.version.Version("2.0.0") _BQ_STORAGE_OPTIONAL_READ_SESSION_VERSION = packaging.version.Version("2.6.0") @@ -247,3 +245,51 @@ def try_import(self, raise_if_error: bool = False) -> Any: and PYARROW_VERSIONS.try_import() is not None and PYARROW_VERSIONS.installed_version >= _MIN_PYARROW_VERSION_RANGE ) + + +class PandasGBQVersions: + """Version and delegation comparisons for pandas-gbq package.""" + + def __init__(self): + self._installed_version = None + self._delegation_api_version = None + + @property + def installed_version(self) -> packaging.version.Version: + """Return the parsed version of pandas-gbq.""" + if self._installed_version is not None: + return self._installed_version + + try: + import pandas_gbq # type: ignore + + self._installed_version = packaging.version.parse( + getattr(pandas_gbq, "__version__", "0.0.0") + ) + except Exception: + self._installed_version = packaging.version.parse("0.0.0") + return self._installed_version + + @property + def delegation_api_version(self) -> int: + """Return the delegation API version of pandas-gbq if installed, otherwise 0.""" + if self._delegation_api_version is not None: + return self._delegation_api_version + + try: + import pandas_gbq # type: ignore + + self._delegation_api_version = int( + getattr(pandas_gbq, "_internal_delegation_api_version", 0) + ) + except Exception: + self._delegation_api_version = 0 + return self._delegation_api_version + + @property + def is_delegation_supported(self) -> bool: + """True if the installed pandas-gbq version supports query delegation API (version >= 1).""" + return self.delegation_api_version >= 1 + + +PANDAS_GBQ_VERSIONS = PandasGBQVersions() diff --git a/packages/google-cloud-bigquery/google/cloud/bigquery/table.py b/packages/google-cloud-bigquery/google/cloud/bigquery/table.py index 870cdcc5d2ab..89ca0f4eec84 100644 --- a/packages/google-cloud-bigquery/google/cloud/bigquery/table.py +++ b/packages/google-cloud-bigquery/google/cloud/bigquery/table.py @@ -21,8 +21,7 @@ import functools import operator import typing -from typing import Any, Dict, Iterable, Iterator, List, Optional, Tuple, Union, Sequence - +from typing import Any, Dict, Iterable, Iterator, List, Optional, Sequence, Tuple, Union import warnings try: @@ -57,29 +56,34 @@ import google.api_core.exceptions from google.api_core.page_iterator import HTTPIterator - import google.cloud._helpers # type: ignore -from google.cloud.bigquery import _helpers -from google.cloud.bigquery import _pandas_helpers -from google.cloud.bigquery import _versions_helpers + +from google.cloud.bigquery import ( + _helpers, + _pandas_helpers, + _string_references, + _versions_helpers, +) from google.cloud.bigquery import exceptions as bq_exceptions +from google.cloud.bigquery import external_config +from google.cloud.bigquery import schema as _schema from google.cloud.bigquery._tqdm_helpers import get_progress_bar from google.cloud.bigquery.encryption_configuration import EncryptionConfiguration from google.cloud.bigquery.enums import DefaultPandasDTypes from google.cloud.bigquery.external_config import ExternalConfig -from google.cloud.bigquery import schema as _schema -from google.cloud.bigquery.schema import _build_schema_resource -from google.cloud.bigquery.schema import _parse_schema_resource -from google.cloud.bigquery.schema import _to_schema_fields -from google.cloud.bigquery import external_config -from google.cloud.bigquery import _string_references +from google.cloud.bigquery.schema import ( + _build_schema_resource, + _parse_schema_resource, + _to_schema_fields, +) if typing.TYPE_CHECKING: # pragma: NO COVER # Unconditionally import optional dependencies again to tell pytype that # they are not None, avoiding false "no attribute" errors. + import geopandas # type: ignore import pandas import pyarrow - import geopandas # type: ignore + from google.cloud import bigquery_storage # type: ignore from google.cloud.bigquery.dataset import DatasetReference @@ -797,7 +801,7 @@ def time_partitioning(self, value): api_repr = value.to_api_repr() elif value is not None: raise ValueError( - "value must be google.cloud.bigquery.table.TimePartitioning " "or None" + "value must be google.cloud.bigquery.table.TimePartitioning or None" ) self._properties[self._PROPERTY_TO_API_FIELD["time_partitioning"]] = api_repr @@ -2801,6 +2805,16 @@ def to_dataframe( create_bqstorage_client = False bqstorage_client = None + if _versions_helpers.PANDAS_GBQ_VERSIONS.is_delegation_supported: + client_info = getattr( + getattr(self.client, "_connection", None), "_client_info", None + ) + if client_info: + ua = client_info.user_agent or "" + if "pandas-gbq" not in ua: + version = _versions_helpers.PANDAS_GBQ_VERSIONS.installed_version + client_info.user_agent = f"{ua} pandas-gbq/{version}".strip() + record_batch = self.to_arrow( progress_bar_type=progress_bar_type, bqstorage_client=bqstorage_client, @@ -2990,8 +3004,7 @@ def to_geodataframe( ) if not geography_columns: raise TypeError( - "There must be at least one GEOGRAPHY column" - " to create a GeoDataFrame" + "There must be at least one GEOGRAPHY column to create a GeoDataFrame" ) if geography_column: diff --git a/packages/google-cloud-bigquery/tests/unit/test__versions_helpers.py b/packages/google-cloud-bigquery/tests/unit/test__versions_helpers.py index 8379c87c18e0..16a0ec50baea 100644 --- a/packages/google-cloud-bigquery/tests/unit/test__versions_helpers.py +++ b/packages/google-cloud-bigquery/tests/unit/test__versions_helpers.py @@ -12,6 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. +import sys from unittest import mock import pytest @@ -31,8 +32,8 @@ except ImportError: pandas = None -from google.cloud.bigquery import _versions_helpers -from google.cloud.bigquery import exceptions +from google import cloud +from google.cloud.bigquery import _versions_helpers, exceptions @pytest.mark.skipif(pyarrow is None, reason="pyarrow is not installed") @@ -59,14 +60,12 @@ def test_try_import_raises_error_w_legacy_pyarrow(): versions.try_import(raise_if_error=True) -@pytest.mark.skipif( - pyarrow is not None, - reason="pyarrow is installed, but this test needs it not to be", -) def test_try_import_raises_error_w_no_pyarrow(): versions = _versions_helpers.PyarrowVersions() - with pytest.raises(exceptions.LegacyPyarrowError): - versions.try_import(raise_if_error=True) + with mock.patch.dict(sys.modules, {"pyarrow": None}): + assert versions.try_import(raise_if_error=False) is None + with pytest.raises(exceptions.LegacyPyarrowError): + versions.try_import(raise_if_error=True) @pytest.mark.skipif(pyarrow is None, reason="pyarrow is not installed") @@ -122,17 +121,21 @@ def test_returns_none_with_legacy_bqstorage(): assert bq_storage is None -@pytest.mark.skipif( - bigquery_storage is not None, - reason="Tests behavior when `google-cloud-bigquery-storage` isn't installed", -) def test_returns_none_with_bqstorage_uninstalled(): - try: - bqstorage_versions = _versions_helpers.BQStorageVersions() - bq_storage = bqstorage_versions.try_import() - except exceptions.LegacyBigQueryStorageError: # pragma: NO COVER - raise ("NotFound error raised when raise_if_error == False.") - assert bq_storage is None + versions = _versions_helpers.BQStorageVersions() + with mock.patch.dict(sys.modules, {"google.cloud.bigquery_storage": None}): + with mock.patch.dict(cloud.__dict__): + cloud.__dict__.pop("bigquery_storage", None) + assert versions.try_import() is None + + +def test_raises_error_with_bqstorage_uninstalled(): + versions = _versions_helpers.BQStorageVersions() + with mock.patch.dict(sys.modules, {"google.cloud.bigquery_storage": None}): + with mock.patch.dict(cloud.__dict__): + cloud.__dict__.pop("bigquery_storage", None) + with pytest.raises(exceptions.BigQueryStorageNotFoundError): + versions.try_import(raise_if_error=True) @pytest.mark.skipif( @@ -220,14 +223,12 @@ def test_try_import_raises_error_w_legacy_pandas(): versions.try_import(raise_if_error=True) -@pytest.mark.skipif( - pandas is not None, - reason="pandas is installed, but this test needs it not to be", -) def test_try_import_raises_error_w_no_pandas(): versions = _versions_helpers.PandasVersions() - with pytest.raises(exceptions.LegacyPandasError): - versions.try_import(raise_if_error=True) + with mock.patch.dict(sys.modules, {"pandas": None}): + assert versions.try_import(raise_if_error=False) is None + with pytest.raises(exceptions.LegacyPandasError): + versions.try_import(raise_if_error=True) @pytest.mark.skipif(pandas is None, reason="pandas is not installed") @@ -246,3 +247,95 @@ def test_installed_pandas_version_returns_parsed_version(): assert version.major == 1 assert version.minor == 1 assert version.micro == 0 + + +def test_installed_pandas_gbq_version_returns_cached(): + versions = _versions_helpers.PandasGBQVersions() + versions._installed_version = object() + assert versions.installed_version is versions._installed_version + + +def test_installed_pandas_gbq_version_returns_parsed_version(): + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.__version__ = "1.2.3" + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + version = versions.installed_version + + assert version.major == 1 + assert version.minor == 2 + assert version.micro == 3 + + +def test_installed_pandas_gbq_version_falls_back_on_import_error(): + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": None}): + version = versions.installed_version + + assert version.major == 0 + assert version.minor == 0 + assert version.micro == 0 + + +def test_installed_pandas_gbq_version_falls_back_on_other_error(): + class CorruptPandasGBQ: + @property + def __version__(self): + raise TypeError("Corrupted package") + + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": CorruptPandasGBQ()}): + version = versions.installed_version + + assert version.major == 0 + assert version.minor == 0 + assert version.micro == 0 + + +def test_pandas_gbq_delegation_api_version_returns_cached(): + versions = _versions_helpers.PandasGBQVersions() + versions._delegation_api_version = object() + assert versions.delegation_api_version is versions._delegation_api_version + + +def test_pandas_gbq_delegation_api_version_returns_value(): + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq._internal_delegation_api_version = 42 + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + version = versions.delegation_api_version + + assert version == 42 + + +def test_pandas_gbq_delegation_api_version_falls_back_on_import_error(): + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": None}): + version = versions.delegation_api_version + + assert version == 0 + + +def test_pandas_gbq_delegation_api_version_falls_back_on_other_error(): + class CorruptPandasGBQ: + @property + def _internal_delegation_api_version(self): + raise TypeError("Corrupted package") + + versions = _versions_helpers.PandasGBQVersions() + with mock.patch.dict(sys.modules, {"pandas_gbq": CorruptPandasGBQ()}): + version = versions.delegation_api_version + + assert version == 0 + + +def test_pandas_gbq_is_delegation_supported_true(): + versions = _versions_helpers.PandasGBQVersions() + versions._delegation_api_version = 1 + assert versions.is_delegation_supported is True + + +def test_pandas_gbq_is_delegation_supported_false(): + versions = _versions_helpers.PandasGBQVersions() + versions._delegation_api_version = 0 + assert versions.is_delegation_supported is False diff --git a/packages/google-cloud-bigquery/tests/unit/test_table.py b/packages/google-cloud-bigquery/tests/unit/test_table.py index 5701143a62d4..baf28d642f66 100644 --- a/packages/google-cloud-bigquery/tests/unit/test_table.py +++ b/packages/google-cloud-bigquery/tests/unit/test_table.py @@ -16,6 +16,7 @@ import datetime import logging import re +import sys import time import types import unittest @@ -5723,6 +5724,217 @@ def test_rowiterator_to_geodataframe_delegation(self, to_dataframe): self.assertEqual([v.__class__.__name__ for v in df.g], ["Point"]) + def test_to_dataframe_delegated_updates_user_agent(self): + pytest.importorskip("db_dtypes") + pandas = pytest.importorskip("pandas") + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.__version__ = "1.0.0" + + mock_client_info = mock.Mock() + mock_client_info.user_agent = "gl-python/3.10.0" + + mock_client = _mock_client() + mock_client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=True, + ): + with mock.patch( + "google.cloud.bigquery._versions_helpers.SUPPORTS_RANGE_PYARROW", + False, + ): + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + row_iterator = self._make_one_from_data( + (("name", "STRING"),), (("foo",),) + ) + row_iterator.client = mock_client + df = row_iterator.to_dataframe( + progress_bar_type="tqdm", timeout=5.0 + ) + self.assertIsInstance(df, pandas.DataFrame) + self.assertEqual( + mock_client_info.user_agent, + "gl-python/3.10.0 pandas-gbq/1.0.0", + ) + + def test_to_dataframe_delegated_does_not_duplicate_user_agent(self): + pytest.importorskip("db_dtypes") + pandas = pytest.importorskip("pandas") + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.__version__ = "1.0.0" + + mock_client_info = mock.Mock() + mock_client_info.user_agent = "gl-python/3.10.0 pandas-gbq/1.0.0" + + mock_client = _mock_client() + mock_client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=True, + ): + with mock.patch( + "google.cloud.bigquery._versions_helpers.SUPPORTS_RANGE_PYARROW", + False, + ): + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + row_iterator = self._make_one_from_data( + (("name", "STRING"),), (("foo",),) + ) + row_iterator.client = mock_client + df = row_iterator.to_dataframe( + progress_bar_type="tqdm", timeout=5.0 + ) + self.assertIsInstance(df, pandas.DataFrame) + self.assertEqual( + mock_client_info.user_agent, + "gl-python/3.10.0 pandas-gbq/1.0.0", + ) + + def test_to_dataframe_delegated_when_client_info_is_none(self): + pytest.importorskip("db_dtypes") + pandas = pytest.importorskip("pandas") + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.__version__ = "1.0.0" + + mock_client = _mock_client() + mock_client._connection = mock.Mock(_client_info=None) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=True, + ): + with mock.patch( + "google.cloud.bigquery._versions_helpers.SUPPORTS_RANGE_PYARROW", + False, + ): + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + row_iterator = self._make_one_from_data( + (("name", "STRING"),), (("foo",),) + ) + row_iterator.client = mock_client + df = row_iterator.to_dataframe( + progress_bar_type="tqdm", timeout=5.0 + ) + self.assertIsInstance(df, pandas.DataFrame) + + def test_to_dataframe_delegated_when_user_agent_is_none(self): + pytest.importorskip("db_dtypes") + pandas = pytest.importorskip("pandas") + mock_pandas_gbq = mock.Mock() + mock_pandas_gbq.__version__ = "1.0.0" + + mock_client_info = mock.Mock() + mock_client_info.user_agent = None + + mock_client = _mock_client() + mock_client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=True, + ): + with mock.patch( + "google.cloud.bigquery._versions_helpers.SUPPORTS_RANGE_PYARROW", + False, + ): + with mock.patch.dict(sys.modules, {"pandas_gbq": mock_pandas_gbq}): + row_iterator = self._make_one_from_data( + (("name", "STRING"),), (("foo",),) + ) + row_iterator.client = mock_client + df = row_iterator.to_dataframe( + progress_bar_type="tqdm", timeout=5.0 + ) + self.assertIsInstance(df, pandas.DataFrame) + self.assertEqual( + mock_client_info.user_agent, + "pandas-gbq/1.0.0", + ) + + def test_to_dataframe_delegated_false_does_not_update_user_agent(self): + pytest.importorskip("db_dtypes") + pandas = pytest.importorskip("pandas") + + mock_client_info = mock.Mock() + mock_client_info.user_agent = "gl-python/3.10.0" + + mock_client = _mock_client() + mock_client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=False, + ): + with mock.patch( + "google.cloud.bigquery._versions_helpers.SUPPORTS_RANGE_PYARROW", + False, + ): + row_iterator = self._make_one_from_data( + (("name", "STRING"),), (("foo",),) + ) + row_iterator.client = mock_client + df = row_iterator.to_dataframe(progress_bar_type="tqdm", timeout=5.0) + self.assertIsInstance(df, pandas.DataFrame) + self.assertEqual( + mock_client_info.user_agent, + "gl-python/3.10.0", + ) + + def test_to_geodataframe_updates_user_agent(self): + pytest.importorskip("geopandas") + row_iterator = self._make_one_from_data( + (("name", "STRING"), ("geog", "GEOGRAPHY")), + (("foo", "Point(0 0)"),), + ) + mock_client_info = mock.Mock(user_agent="test-agent") + row_iterator.client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=True, + ): + with mock.patch.object(row_iterator, "to_arrow") as mock_to_arrow: + import pyarrow as pa + + mock_to_arrow.return_value = pa.RecordBatch.from_arrays( + [pa.array(["foo"]), pa.array(["Point(0 0)"])], + names=["name", "geog"], + ) + _ = row_iterator.to_geodataframe(create_bqstorage_client=False) + self.assertIn("pandas-gbq/", mock_client_info.user_agent) + + def test_to_geodataframe_delegated_false_does_not_update_user_agent(self): + pytest.importorskip("geopandas") + row_iterator = self._make_one_from_data( + (("name", "STRING"), ("geog", "GEOGRAPHY")), + (("foo", "Point(0 0)"),), + ) + mock_client_info = mock.Mock(user_agent="test-agent") + row_iterator.client._connection = mock.Mock(_client_info=mock_client_info) + + with mock.patch( + "google.cloud.bigquery._versions_helpers.PandasGBQVersions.is_delegation_supported", + new_callable=mock.PropertyMock, + return_value=False, + ): + with mock.patch.object(row_iterator, "to_arrow") as mock_to_arrow: + import pyarrow as pa + + mock_to_arrow.return_value = pa.RecordBatch.from_arrays( + [pa.array(["foo"]), pa.array(["Point(0 0)"])], + names=["name", "geog"], + ) + _ = row_iterator.to_geodataframe(create_bqstorage_client=False) + self.assertEqual(mock_client_info.user_agent, "test-agent") + class TestPartitionRange(unittest.TestCase): def _get_target_class(self):