Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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()
45 changes: 29 additions & 16 deletions packages/google-cloud-bigquery/google/cloud/bigquery/table.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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:
Expand Down
141 changes: 117 additions & 24 deletions packages/google-cloud-bigquery/tests/unit/test__versions_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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")
Expand All @@ -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")
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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")
Expand All @@ -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
Loading
Loading