From 7d5d1bdbfa8069dd261492457c40563fae3cd393 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Wed, 28 Feb 2024 01:25:43 +0000 Subject: [PATCH 1/9] feat: add configuration option to read_gbq --- bigframes/pandas/__init__.py | 2 + bigframes/session/__init__.py | 67 +++++++++++++++++-- tests/system/small/test_session.py | 28 ++++++++ .../bigframes_vendored/pandas/io/gbq.py | 6 ++ 4 files changed, 98 insertions(+), 5 deletions(-) diff --git a/bigframes/pandas/__init__.py b/bigframes/pandas/__init__.py index 110978a7f1..1da22ce050 100644 --- a/bigframes/pandas/__init__.py +++ b/bigframes/pandas/__init__.py @@ -490,6 +490,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: vendored_pandas_gbq.FiltersType = (), use_cache: bool = True, @@ -501,6 +502,7 @@ def read_gbq( query_or_table, index_col=index_col, columns=columns, + configuration=configuration, max_results=max_results, filters=filters, use_cache=use_cache, diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index 20dd39c0fa..1e6a19ff09 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -16,6 +16,7 @@ from __future__ import annotations +import copy import datetime import itertools import logging @@ -244,9 +245,10 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: third_party_pandas_gbq.FiltersType = (), - use_cache: bool = True, + use_cache: Optional[bool] = None, col_order: Iterable[str] = (), # Add a verify index argument that fails if the index is not unique. ) -> dataframe.DataFrame: @@ -267,6 +269,7 @@ def read_gbq( query_or_table, index_col=index_col, columns=columns, + configuration=configuration, max_results=max_results, api_name="read_gbq", use_cache=use_cache, @@ -275,13 +278,20 @@ def read_gbq( # TODO(swast): Query the snapshot table but mark it as a # deterministic query so we can avoid serializing if we have a # unique index. + if configuration: + raise ValueError( + "The 'configuration' argument is not allowed when " + "directly reading from a table. Please remove " + "'configuration' or use a query." + ) + return self._read_gbq_table( query_or_table, index_col=index_col, columns=columns, max_results=max_results, api_name="read_gbq", - use_cache=use_cache, + use_cache=use_cache if use_cache is not None else True, ) def _to_query( @@ -365,6 +375,7 @@ def _query_to_destination( query: str, index_cols: List[str], api_name: str, + configuration: Optional[dict] = None, use_cache: bool = True, ) -> Tuple[Optional[bigquery.TableReference], Optional[bigquery.QueryJob]]: # If a dry_run indicates this is not a query type job, then don't @@ -387,7 +398,11 @@ def _query_to_destination( ][:_MAX_CLUSTER_COLUMNS] temp_table = self._create_empty_temp_table(schema, cluster_cols) - job_config = bigquery.QueryJobConfig() + job_config = ( + bigquery.QueryJobConfig.from_api_repr(configuration) + if configuration + else bigquery.QueryJobConfig() + ) job_config.labels["bigframes-api"] = api_name job_config.destination = temp_table job_config.use_query_cache = use_cache @@ -412,8 +427,9 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: Optional[dict] = None, max_results: Optional[int] = None, - use_cache: bool = True, + use_cache: Optional[bool] = None, col_order: Iterable[str] = (), ) -> dataframe.DataFrame: """Turn a SQL query into a DataFrame. @@ -477,6 +493,7 @@ def read_gbq_query( query=query, index_col=index_col, columns=columns, + configuration=configuration, max_results=max_results, api_name="read_gbq_query", use_cache=use_cache, @@ -488,10 +505,27 @@ def _read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: Optional[dict] = None, max_results: Optional[int] = None, api_name: str = "read_gbq_query", - use_cache: bool = True, + use_cache: Optional[bool] = None, ) -> dataframe.DataFrame: + configuration = _transform_read_gbq_configuration(configuration) + if configuration and "query" in configuration: + if "query" in configuration["query"]: + raise ValueError( + "The query statement must not be included in the ", + "'configuration' because it is already provided as", + " a separate parameter.", + ) + if ("useQueryCache" in configuration["query"]) and (use_cache is not None): + raise ValueError( + "'useQueryCache' in 'configuration' conflicts with" + " 'use_cache' parameter. Please specify only one." + ) + + use_cache = use_cache if use_cache is not None else True + if isinstance(index_col, str): index_cols = [index_col] else: @@ -501,6 +535,7 @@ def _read_gbq_query( query, index_cols, api_name=api_name, + configuration=configuration, use_cache=use_cache, ) @@ -1722,3 +1757,25 @@ def _convert_to_nonnull_string(column: ibis_types.Column) -> ibis_types.StringVa # Escape backslashes and use backslash as delineator escaped = typing.cast(ibis_types.StringColumn, result.fillna("")).replace("\\", "\\\\") # type: ignore return typing.cast(ibis_types.StringColumn, ibis.literal("\\")).concat(escaped) + + +def _transform_read_gbq_configuration(configuration): + """ + For backwards-compatibility, convert any previously client-side only + parameters such as timeoutMs to the property name expected by the REST API. + + Makes a copy of configuration if changes are needed. + """ + + if configuration is None: + return None + + timeout_ms = configuration.get("query", {}).get("timeoutMs") + if timeout_ms is not None: + # Transform timeoutMs to an actual server-side configuration. + # https://github.com/googleapis/python-bigquery-pandas/issues/479 + configuration = copy.deepcopy(configuration) + del configuration["query"]["timeoutMs"] + configuration["jobTimeoutMs"] = timeout_ms + + return configuration diff --git a/tests/system/small/test_session.py b/tests/system/small/test_session.py index 85573472b9..3b8e2cd013 100644 --- a/tests/system/small/test_session.py +++ b/tests/system/small/test_session.py @@ -20,6 +20,7 @@ import typing from typing import List +from google.api_core.exceptions import InternalServerError import google.cloud.bigquery as bigquery import numpy as np import pandas as pd @@ -353,6 +354,33 @@ def test_read_gbq_table_wildcard_with_filter(session: bigframes.Session): assert df.shape == (348485, 32) +@pytest.mark.parametrize( + ("config"), + [ + {"query": {"useQueryCache": True, "maximumBytesBilled": "1000000000"}}, + pytest.param( + {"query": {"useQueryCache": False, "maximumBytesBilled": "100"}}, + marks=pytest.mark.xfail( + raises=InternalServerError, + ), + ), + ], +) +def test_read_gbq_with_configuration( + session: bigframes.Session, scalars_table_id: str, config: dict +): + query = f"""SELECT + t.float64_col * 2 AS my_floats, + CONCAT(t.string_col, "_2") AS my_strings, + t.int64_col > 0 AS my_bools, + FROM `{scalars_table_id}` AS t + """ + + df = session.read_gbq(query, configuration=config) + + assert df.shape == (9, 3) + + def test_read_gbq_model(session, penguins_linear_model_name): model = session.read_gbq_model(penguins_linear_model_name) assert isinstance(model, bigframes.ml.linear_model.LinearRegression) diff --git a/third_party/bigframes_vendored/pandas/io/gbq.py b/third_party/bigframes_vendored/pandas/io/gbq.py index 1f31c530d2..cdaa294688 100644 --- a/third_party/bigframes_vendored/pandas/io/gbq.py +++ b/third_party/bigframes_vendored/pandas/io/gbq.py @@ -19,6 +19,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: FiltersType = (), use_cache: bool = True, @@ -107,6 +108,11 @@ def read_gbq( columns (Iterable[str]): List of BigQuery column names in the desired order for results DataFrame. + configuration (dict, optional): + Query config parameters for job processing. + For example: configuration = {'query': {'useQueryCache': False}}. + For more information see `BigQuery REST API Reference + `__. max_results (Optional[int], default None): If set, limit the maximum number of rows to fetch from the query results. From 4bea5128f0098b659462cd2983d20435549fa30c Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Wed, 28 Feb 2024 20:18:14 +0000 Subject: [PATCH 2/9] Update logic. --- bigframes/session/__init__.py | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index 1e6a19ff09..4baa26580d 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -376,7 +376,7 @@ def _query_to_destination( index_cols: List[str], api_name: str, configuration: Optional[dict] = None, - use_cache: bool = True, + use_cache: Optional[bool] = None, ) -> Tuple[Optional[bigquery.TableReference], Optional[bigquery.QueryJob]]: # If a dry_run indicates this is not a query type job, then don't # bother trying to do a CREATE TEMP TABLE ... AS SELECT ... statement. @@ -405,7 +405,8 @@ def _query_to_destination( ) job_config.labels["bigframes-api"] = api_name job_config.destination = temp_table - job_config.use_query_cache = use_cache + if not job_config.use_query_cache: + job_config.use_query_cache = use_cache if use_cache is not None else True try: # Write to temp table to workaround BigQuery 10 GB query results @@ -523,9 +524,6 @@ def _read_gbq_query( "'useQueryCache' in 'configuration' conflicts with" " 'use_cache' parameter. Please specify only one." ) - - use_cache = use_cache if use_cache is not None else True - if isinstance(index_col, str): index_cols = [index_col] else: From 9589f9825d0cd3b1a22df90ea9f1d92e6049898e Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Wed, 13 Mar 2024 23:32:21 +0000 Subject: [PATCH 3/9] Update logic and add new test. --- bigframes/pandas/__init__.py | 8 ++- bigframes/session/__init__.py | 61 ++++++++++++------- tests/system/small/test_session.py | 15 ++++- .../bigframes_vendored/pandas/io/gbq.py | 10 +-- 4 files changed, 63 insertions(+), 31 deletions(-) diff --git a/bigframes/pandas/__init__.py b/bigframes/pandas/__init__.py index 1da22ce050..57f5853d50 100644 --- a/bigframes/pandas/__init__.py +++ b/bigframes/pandas/__init__.py @@ -490,10 +490,10 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: dict = {}, max_results: Optional[int] = None, filters: vendored_pandas_gbq.FiltersType = (), - use_cache: bool = True, + use_cache: Optional[bool] = None, col_order: Iterable[str] = (), ) -> bigframes.dataframe.DataFrame: _set_default_session_location_if_possible(query_or_table) @@ -528,8 +528,9 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), + configuration: dict = {}, max_results: Optional[int] = None, - use_cache: bool = True, + use_cache: Optional[bool] = None, col_order: Iterable[str] = (), ) -> bigframes.dataframe.DataFrame: _set_default_session_location_if_possible(query) @@ -538,6 +539,7 @@ def read_gbq_query( query, index_col=index_col, columns=columns, + configuration=configuration, max_results=max_results, use_cache=use_cache, col_order=col_order, diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index 4baa26580d..b7d7de74a2 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -245,7 +245,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: dict = {}, max_results: Optional[int] = None, filters: third_party_pandas_gbq.FiltersType = (), use_cache: Optional[bool] = None, @@ -375,8 +375,7 @@ def _query_to_destination( query: str, index_cols: List[str], api_name: str, - configuration: Optional[dict] = None, - use_cache: Optional[bool] = None, + configuration: dict = {"query": {"useQueryCache": True}}, ) -> Tuple[Optional[bigquery.TableReference], Optional[bigquery.QueryJob]]: # If a dry_run indicates this is not a query type job, then don't # bother trying to do a CREATE TEMP TABLE ... AS SELECT ... statement. @@ -398,28 +397,35 @@ def _query_to_destination( ][:_MAX_CLUSTER_COLUMNS] temp_table = self._create_empty_temp_table(schema, cluster_cols) - job_config = ( - bigquery.QueryJobConfig.from_api_repr(configuration) - if configuration - else bigquery.QueryJobConfig() + timeout_ms = configuration.get("jobTimeoutMs") or configuration["query"].get( + "timeoutMs" + ) + + # Convert timeout_ms to seconds, ensuring a minimum of 0.1 seconds to avoid + # the program getting stuck on too-short timeouts. + timeout = max(int(timeout_ms) * 1e-3, 0.1) if timeout_ms else None + + job_config = typing.cast( + bigquery.QueryJobConfig, + bigquery.QueryJobConfig.from_api_repr(configuration), ) job_config.labels["bigframes-api"] = api_name job_config.destination = temp_table - if not job_config.use_query_cache: - job_config.use_query_cache = use_cache if use_cache is not None else True try: # Write to temp table to workaround BigQuery 10 GB query results # limit. See: internal issue 303057336. job_config.labels["error_caught"] = "true" - _, query_job = self._start_query(query, job_config=job_config) + _, query_job = self._start_query( + query, job_config=job_config, timeout=timeout + ) return query_job.destination, query_job except google.api_core.exceptions.BadRequest: # Some SELECT statements still aren't compatible with cluster # tables as the destination. For example, if the query has a # top-level ORDER BY, this conflicts with our ability to cluster # the table by the index column(s). - _, query_job = self._start_query(query) + _, query_job = self._start_query(query, timeout=timeout) return query_job.destination, query_job def read_gbq_query( @@ -428,7 +434,7 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: dict = {}, max_results: Optional[int] = None, use_cache: Optional[bool] = None, col_order: Iterable[str] = (), @@ -506,24 +512,34 @@ def _read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: dict = {}, max_results: Optional[int] = None, api_name: str = "read_gbq_query", use_cache: Optional[bool] = None, ) -> dataframe.DataFrame: configuration = _transform_read_gbq_configuration(configuration) - if configuration and "query" in configuration: - if "query" in configuration["query"]: - raise ValueError( - "The query statement must not be included in the ", - "'configuration' because it is already provided as", - " a separate parameter.", - ) - if ("useQueryCache" in configuration["query"]) and (use_cache is not None): + + if not "query" in configuration: + configuration["query"] = {} + + if "query" in configuration["query"]: + raise ValueError( + "The query statement must not be included in the ", + "'configuration' because it is already provided as", + " a separate parameter.", + ) + + if "useQueryCache" in configuration["query"]: + if use_cache is not None: raise ValueError( "'useQueryCache' in 'configuration' conflicts with" " 'use_cache' parameter. Please specify only one." ) + else: + configuration["query"]["useQueryCache"] = ( + True if use_cache is None else use_cache + ) + if isinstance(index_col, str): index_cols = [index_col] else: @@ -534,7 +550,6 @@ def _read_gbq_query( index_cols, api_name=api_name, configuration=configuration, - use_cache=use_cache, ) # If there was no destination table, that means the query must have @@ -558,7 +573,7 @@ def _read_gbq_query( index_col=index_cols, columns=columns, max_results=max_results, - use_cache=use_cache, + use_cache=configuration["query"]["useQueryCache"], ) def read_gbq_table( diff --git a/tests/system/small/test_session.py b/tests/system/small/test_session.py index 3b8e2cd013..3583ac3f54 100644 --- a/tests/system/small/test_session.py +++ b/tests/system/small/test_session.py @@ -20,6 +20,7 @@ import typing from typing import List +import google from google.api_core.exceptions import InternalServerError import google.cloud.bigquery as bigquery import numpy as np @@ -357,7 +358,19 @@ def test_read_gbq_table_wildcard_with_filter(session: bigframes.Session): @pytest.mark.parametrize( ("config"), [ - {"query": {"useQueryCache": True, "maximumBytesBilled": "1000000000"}}, + { + "query": { + "useQueryCache": True, + "maximumBytesBilled": "1000000000", + "timeoutMs": 10000, + } + }, + pytest.param( + {"query": {"useQueryCache": True, "timeoutMs": 50}}, + marks=pytest.mark.xfail( + raises=google.api_core.exceptions.BadRequest, + ), + ), pytest.param( {"query": {"useQueryCache": False, "maximumBytesBilled": "100"}}, marks=pytest.mark.xfail( diff --git a/third_party/bigframes_vendored/pandas/io/gbq.py b/third_party/bigframes_vendored/pandas/io/gbq.py index cdaa294688..1a6f656def 100644 --- a/third_party/bigframes_vendored/pandas/io/gbq.py +++ b/third_party/bigframes_vendored/pandas/io/gbq.py @@ -19,10 +19,10 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: dict = {}, max_results: Optional[int] = None, filters: FiltersType = (), - use_cache: bool = True, + use_cache: Optional[bool] = None, col_order: Iterable[str] = (), ): """Loads a DataFrame from BigQuery. @@ -127,8 +127,10 @@ def read_gbq( If using wildcard table suffix in query_or_table, can specify '_table_suffix' pseudo column to filter the tables to be read into the DataFrame. - use_cache (bool, default True): - Whether to cache the query inputs. Default to True. + use_cache (Optional[bool], default None): + Caches query results if set to `True`. When `None`, it behaves + as `True`, but should not be combined with `useQueryCache` in + `configuration` to avoid conflicts. col_order (Iterable[str]): Alias for columns, retained for backwards compatibility. From a312395033a680b2b496ff4951656cede3736309 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Wed, 13 Mar 2024 23:50:29 +0000 Subject: [PATCH 4/9] Update code logic for timeout after merge. --- bigframes/session/__init__.py | 3 ++- bigframes/session/_io/bigquery.py | 3 ++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index 07f64a4e20..96e5ff1b71 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -1631,13 +1631,14 @@ def _start_query( sql: str, job_config: Optional[bigquery.job.QueryJobConfig] = None, max_results: Optional[int] = None, + timeout: Optional[float] = None, ) -> Tuple[bigquery.table.RowIterator, bigquery.QueryJob]: """ Starts BigQuery query job and waits for results. """ job_config = self._prepare_query_job_config(job_config) return bigframes.session._io.bigquery.start_query_with_client( - self.bqclient, sql, job_config, max_results + self.bqclient, sql, job_config, max_results, timeout ) def _start_query_create_model( diff --git a/bigframes/session/_io/bigquery.py b/bigframes/session/_io/bigquery.py index 67820bbbcb..38ff7429ec 100644 --- a/bigframes/session/_io/bigquery.py +++ b/bigframes/session/_io/bigquery.py @@ -220,6 +220,7 @@ def start_query_with_client( sql: str, job_config: bigquery.job.QueryJobConfig, max_results: Optional[int] = None, + timeout: Optional[float] = None, ) -> Tuple[bigquery.table.RowIterator, bigquery.QueryJob]: """ Starts query job and waits for results. @@ -230,7 +231,7 @@ def start_query_with_client( ) try: - query_job = bq_client.query(sql, job_config=job_config) + query_job = bq_client.query(sql, job_config=job_config, timeout=timeout) except google.api_core.exceptions.Forbidden as ex: if "Drive credentials" in ex.message: ex.message += "\nCheck https://cloud.google.com/bigquery/docs/query-drive-data#Google_Drive_permissions." From 3d3e1166933dc0c42e664d2445060f58524af4f3 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Wed, 13 Mar 2024 23:55:43 +0000 Subject: [PATCH 5/9] Update lint. --- bigframes/session/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index 96e5ff1b71..c79dc24b3d 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -530,7 +530,7 @@ def _read_gbq_query( ) -> dataframe.DataFrame: configuration = _transform_read_gbq_configuration(configuration) - if not "query" in configuration: + if "query" not in configuration: configuration["query"] = {} if "query" in configuration["query"]: From 34f5903709b6c412e47f12ae031aa07cf2c8bf16 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Thu, 14 Mar 2024 17:34:46 +0000 Subject: [PATCH 6/9] update annotation and default value --- bigframes/pandas/__init__.py | 4 ++-- bigframes/session/__init__.py | 10 +++++----- third_party/bigframes_vendored/pandas/io/gbq.py | 2 +- 3 files changed, 8 insertions(+), 8 deletions(-) diff --git a/bigframes/pandas/__init__.py b/bigframes/pandas/__init__.py index e69e38acc1..e21103115c 100644 --- a/bigframes/pandas/__init__.py +++ b/bigframes/pandas/__init__.py @@ -491,7 +491,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: vendored_pandas_gbq.FiltersType = (), use_cache: Optional[bool] = None, @@ -529,7 +529,7 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, use_cache: Optional[bool] = None, col_order: Iterable[str] = (), diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index c06da6ef88..c600e79914 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -256,7 +256,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: third_party_pandas_gbq.FiltersType = (), use_cache: Optional[bool] = None, @@ -445,7 +445,7 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, use_cache: Optional[bool] = None, col_order: Iterable[str] = (), @@ -523,7 +523,7 @@ def _read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, api_name: str = "read_gbq_query", use_cache: Optional[bool] = None, @@ -1854,7 +1854,7 @@ def _convert_to_nonnull_string(column: ibis_types.Column) -> ibis_types.StringVa return typing.cast(ibis_types.StringColumn, ibis.literal("\\")).concat(escaped) -def _transform_read_gbq_configuration(configuration): +def _transform_read_gbq_configuration(configuration: Optional[dict]) -> dict: """ For backwards-compatibility, convert any previously client-side only parameters such as timeoutMs to the property name expected by the REST API. @@ -1863,7 +1863,7 @@ def _transform_read_gbq_configuration(configuration): """ if configuration is None: - return None + return {} timeout_ms = configuration.get("query", {}).get("timeoutMs") if timeout_ms is not None: diff --git a/third_party/bigframes_vendored/pandas/io/gbq.py b/third_party/bigframes_vendored/pandas/io/gbq.py index 1a6f656def..3ec627d76d 100644 --- a/third_party/bigframes_vendored/pandas/io/gbq.py +++ b/third_party/bigframes_vendored/pandas/io/gbq.py @@ -19,7 +19,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: dict = {}, + configuration: Optional[dict] = None, max_results: Optional[int] = None, filters: FiltersType = (), use_cache: Optional[bool] = None, From ee9eb2cf603cecc5b74694118d784dc9744f62b9 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Thu, 14 Mar 2024 17:36:52 +0000 Subject: [PATCH 7/9] update config check. --- bigframes/session/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index c600e79914..addc6e94dc 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -289,7 +289,7 @@ def read_gbq( # TODO(swast): Query the snapshot table but mark it as a # deterministic query so we can avoid serializing if we have a # unique index. - if configuration: + if configuration is not None: raise ValueError( "The 'configuration' argument is not allowed when " "directly reading from a table. Please remove " From 742b1f1e8dc2097d60aeae20a5b1a272c6a7fb60 Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Thu, 21 Mar 2024 18:06:57 +0000 Subject: [PATCH 8/9] Update test cases --- bigframes/pandas/__init__.py | 4 ++-- bigframes/session/__init__.py | 6 +++--- tests/system/small/test_session.py | 7 +++++-- third_party/bigframes_vendored/pandas/io/gbq.py | 4 ++-- 4 files changed, 12 insertions(+), 9 deletions(-) diff --git a/bigframes/pandas/__init__.py b/bigframes/pandas/__init__.py index 3fd27c9e50..c007927a8a 100644 --- a/bigframes/pandas/__init__.py +++ b/bigframes/pandas/__init__.py @@ -491,7 +491,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, filters: vendored_pandas_gbq.FiltersType = (), use_cache: Optional[bool] = None, @@ -529,7 +529,7 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, use_cache: Optional[bool] = None, col_order: Iterable[str] = (), diff --git a/bigframes/session/__init__.py b/bigframes/session/__init__.py index aecc53fa10..e56535e9a6 100644 --- a/bigframes/session/__init__.py +++ b/bigframes/session/__init__.py @@ -256,7 +256,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, filters: third_party_pandas_gbq.FiltersType = (), use_cache: Optional[bool] = None, @@ -445,7 +445,7 @@ def read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, use_cache: Optional[bool] = None, col_order: Iterable[str] = (), @@ -523,7 +523,7 @@ def _read_gbq_query( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, api_name: str = "read_gbq_query", use_cache: Optional[bool] = None, diff --git a/tests/system/small/test_session.py b/tests/system/small/test_session.py index 8b84a79913..fdd2f44120 100644 --- a/tests/system/small/test_session.py +++ b/tests/system/small/test_session.py @@ -21,7 +21,6 @@ from typing import List import google -from google.api_core.exceptions import InternalServerError import google.cloud.bigquery as bigquery import numpy as np import pandas as pd @@ -369,12 +368,16 @@ def test_read_gbq_table_wildcard_with_filter(session: bigframes.Session): {"query": {"useQueryCache": True, "timeoutMs": 50}}, marks=pytest.mark.xfail( raises=google.api_core.exceptions.BadRequest, + reason="Expected failure due to timeout being set too short.", + match=r"API deadline too short", ), ), pytest.param( {"query": {"useQueryCache": False, "maximumBytesBilled": "100"}}, marks=pytest.mark.xfail( - raises=InternalServerError, + raises=google.api_core.exceptions.InternalServerError, + reason="Expected failure when the query exceeds the maximum bytes billed limit.", + match=r"Query exceeded limit for bytes billed", ), ), ], diff --git a/third_party/bigframes_vendored/pandas/io/gbq.py b/third_party/bigframes_vendored/pandas/io/gbq.py index 3ec627d76d..050fcf8a59 100644 --- a/third_party/bigframes_vendored/pandas/io/gbq.py +++ b/third_party/bigframes_vendored/pandas/io/gbq.py @@ -3,7 +3,7 @@ from __future__ import annotations -from typing import Any, Iterable, Literal, Optional, Tuple, Union +from typing import Any, Dict, Iterable, Literal, Optional, Tuple, Union from bigframes import constants @@ -19,7 +19,7 @@ def read_gbq( *, index_col: Iterable[str] | str = (), columns: Iterable[str] = (), - configuration: Optional[dict] = None, + configuration: Optional[Dict] = None, max_results: Optional[int] = None, filters: FiltersType = (), use_cache: Optional[bool] = None, From 2e80ec9d9cd4a338ca19785eab983a018023b80c Mon Sep 17 00:00:00 2001 From: Huan Chen Date: Thu, 21 Mar 2024 19:33:05 +0000 Subject: [PATCH 9/9] Update test --- tests/system/small/test_session.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/tests/system/small/test_session.py b/tests/system/small/test_session.py index fdd2f44120..f314b205c6 100644 --- a/tests/system/small/test_session.py +++ b/tests/system/small/test_session.py @@ -369,7 +369,6 @@ def test_read_gbq_table_wildcard_with_filter(session: bigframes.Session): marks=pytest.mark.xfail( raises=google.api_core.exceptions.BadRequest, reason="Expected failure due to timeout being set too short.", - match=r"API deadline too short", ), ), pytest.param( @@ -377,7 +376,6 @@ def test_read_gbq_table_wildcard_with_filter(session: bigframes.Session): marks=pytest.mark.xfail( raises=google.api_core.exceptions.InternalServerError, reason="Expected failure when the query exceeds the maximum bytes billed limit.", - match=r"Query exceeded limit for bytes billed", ), ), ],