Repository navigation
Conversation
…hods and add REST fallback Convert RowIterator Arrow streaming static helpers to instance methods. Defer BQ Storage gRPC client initialization until after inline REST decoding and add REST pageToken pagination fallback.
Decode the server-provided REST arrowSchema IPC header in to_arrow() when total_rows == 0, guaranteeing 100% exact Arrow schema type alignment for empty query results instead of falling back to JSON schema conversion.
…sult Remove experimental query_arrow() method and ArrowQueryResult wrapper. RowIterator natively provides job metadata and to_arrow() returning a pyarrow.Table directly, providing a cleaner, more pythonic interface.
Remove custom ArrowSerializationOptions wrapper class since bufferByteLimit and useInt64Timestamp do not exist on the backend ArrowSerializationOptions proto, and compression is already exposed via compression_codec on query_and_wait(). Remove unsupported REST jobs.getQueryResults pageToken Arrow fallback from RowIterator._download_arrow_from_job_id(). Complete Page 1 over REST without requiring google-cloud-bigquery-storage, and stream Page 2+ via BigQuery Storage Read API gRPC default stream. Propagate job_config.query_results_format in _job_helpers.query_and_wait() and update unit tests, system tests, and 3-way benchmark script.
Validate bqstorage_client upfront in _download_arrow_from_job_id() when more_pages_needed is true before yielding the initial REST Arrow batch. Single-page queries remain gRPC-free while multi-page queries fail on the first iteration if google-cloud-bigquery-storage is missing. Also remove obsolete prototype unit tests from test_client.py.
Remove _raw_first_page_response and the special-case 0-row schema fallback in RowIterator.to_arrow() so RowIterator.to_arrow() remains untouched while to_arrow_iterable() handles 0-row queries with zero gRPC calls. Add unit and system tests for 0-row to_arrow_iterable() and multi-page first_page_response preservation in QueryJob.result(), and update benchmark_arrow_vs_json.py to use to_arrow_iterable().
Serialize compression_codec into top-level arrowSerializationOptions on the jobs.query request body instead of nesting it under formatOptions, where the BigQuery REST API ignored it. Add arrowSerializationOptions to the _supported_by_jobs_query allowlist.
Ensure _download_arrow_from_job_id closes bqstorage_client._transport in a finally block when self.client._ensure_bqstorage_client() creates a new BigQueryReadClient instance, preventing gRPC channel and socket leaks while leaving caller-supplied clients untouched.
Short-circuit redundant jobs.get RPCs in query_and_wait when inline REST Arrow batches satisfy max_results. Slice inline and gRPC Arrow record batches to enforce max_results across page boundaries, and preserve server Arrow schemas on zero-row and max_results=0 tables.
Expose arrow_serialization_options on BigQueryReadClient.read_rows and ReadRowsStream so callers reading job streams can configure Arrow buffer compression and timestamp precision while preserving automatic stream resumption on reconnect.
Propagate compression_codec from query_and_wait through RowIterator to read_rows on Page 2+ and jobs.insert streams. Fall back to the GAPIC read_rows request dict when bigquery-storage predates arrow_serialization_options, and handle cached multi-page responses.
Replace direct GapicBigQueryReadClient invocation with a public read_rows fallback when older google-cloud-bigquery-storage releases reject arrow_serialization_options. Emit a UserWarning so callers know multi-page Arrow streams remain uncompressed until upgraded, while preserving automatic stream retry and reconnect behavior.
Disable query caching in test_query_arrow_multi_page_compression_codec so uncompressed and compressed runs both execute uncached. Add test_query_arrow_cached_multi_page to explicitly verify cached multi-page Arrow result pagination against BigQuery.
Add test_read_rows_reconnect_preserves_arrow_serialization_options to confirm ReadRowsStream._reconnect forwards arrow_serialization_options with the updated row offset after transient stream disconnects.
Add query_results_format and compression_codec parameters to the RowIterator class docstring so documented arguments match __init__.
BigQuery Storage Read API streams/_default only supports anonymous query result tables and rejects explicit destination tables with INVALID_ARGUMENT. When query_and_wait falls back to jobs.insert with an explicit destination table, open a single-stream Arrow read session via create_read_session (passing compression_codec in read_options) and stream batches over read_rows while preserving 0-row Arrow schemas.
There was a problem hiding this comment.
Code Review
This pull request adds support for Apache Arrow serialization options and compression codecs (such as LZ4 and ZSTD) across the BigQuery and BigQuery Storage clients, enabling faster multi-page query results and introducing a performance benchmark. Feedback on these changes highlights a critical bug in table.py where a schema-applying elif block is unreachable, code duplication in reader.py when calling read_rows, a security risk from a hardcoded sandbox project ID in the benchmark script, and the use of a private _transport attribute instead of the public close() method for client cleanup.
| or self._query_results_format == QueryResultsFormat.ARROW.value | ||
| ): | ||
| return pyarrow.Table.from_batches(record_batches) | ||
| elif ( | ||
| self._query_results_format == QueryResultsFormat.ARROW.value | ||
| and self._arrow_schema is not None | ||
| ): | ||
| return pyarrow.Table.from_batches(record_batches, schema=self._arrow_schema) |
There was a problem hiding this comment.
CRITICAL BUG: The elif block checking self._arrow_schema is not None is dead code and will never be executed.
Because the preceding if condition checks self._query_results_format == QueryResultsFormat.ARROW.value without checking self._arrow_schema, any time the format is ARROW, the if block will match and return pyarrow.Table.from_batches(record_batches) immediately, completely bypassing the elif block that applies the schema.
To fix this, we should modify the if condition to only match when self._arrow_schema is None.
| or self._query_results_format == QueryResultsFormat.ARROW.value | |
| ): | |
| return pyarrow.Table.from_batches(record_batches) | |
| elif ( | |
| self._query_results_format == QueryResultsFormat.ARROW.value | |
| and self._arrow_schema is not None | |
| ): | |
| return pyarrow.Table.from_batches(record_batches, schema=self._arrow_schema) | |
| or ( | |
| self._query_results_format == QueryResultsFormat.ARROW.value | |
| and self._arrow_schema is None | |
| ) | |
| ): | |
| return pyarrow.Table.from_batches(record_batches) | |
| elif ( | |
| self._query_results_format == QueryResultsFormat.ARROW.value | |
| and self._arrow_schema is not None | |
| ): | |
| return pyarrow.Table.from_batches(record_batches, schema=self._arrow_schema) |
| if self._arrow_serialization_options is not None: | ||
| self._wrapped = self._client.read_rows( | ||
| request={ | ||
| "read_stream": self._name, | ||
| "offset": self._offset, | ||
| "arrow_serialization_options": self._arrow_serialization_options, | ||
| }, | ||
| **self._read_rows_kwargs, | ||
| ) | ||
| else: | ||
| self._wrapped = self._client.read_rows( | ||
| read_stream=self._name, | ||
| offset=self._offset, | ||
| **self._read_rows_kwargs, | ||
| ) |
There was a problem hiding this comment.
The if/else block duplicates the call to self._client.read_rows. We can simplify this by dynamically constructing the keyword arguments dictionary, which improves readability and maintainability.
kwargs = {
"read_stream": self._name,
"offset": self._offset,
**self._read_rows_kwargs,
}
if self._arrow_serialization_options is not None:
kwargs["arrow_serialization_options"] = self._arrow_serialization_options
self._wrapped = self._client.read_rows(**kwargs)| "GOOGLE_APPLICATION_CREDENTIALS", | ||
| os.path.expanduser("~/.config/gcloud/application_default_credentials.json"), | ||
| ) | ||
| os.environ.setdefault("GCLOUD_PROJECT", "omairn-project-sandbox-174760") |
There was a problem hiding this comment.
Avoid committing hardcoded sandbox project IDs or credentials to the repository. This is a security and privacy risk, and makes the benchmark script brittle for other users who do not have access to this project.
Instead, let the user specify the project ID via the GCLOUD_PROJECT environment variable. Ensure you validate the environment variable and fail fast if it is missing, empty, or contains only whitespace, rather than silently falling back to a default value.
| os.environ.setdefault("GCLOUD_PROJECT", "omairn-project-sandbox-174760") | |
| project_id = os.environ.get("GCLOUD_PROJECT") | |
| if not project_id or not project_id.strip(): | |
| raise ValueError("GCLOUD_PROJECT environment variable must be set to a valid, non-empty value.") |
References
- When parsing environment variables, if a user explicitly sets an invalid value (such as whitespace-only), fail fast and raise an error to notify them of the invalid configuration rather than silently falling back to a default value.
| if owns_bqstorage_client and bqstorage_client is not None: | ||
| bqstorage_client._transport.close() |
There was a problem hiding this comment.
Avoid accessing the private _transport attribute to close the client. Instead, use the public close() method on the client itself, which is cleaner and more robust against future changes in google-api-core.
Additionally, when closing resources during cleanup, wrap the close call in a try-except block to handle exceptions individually, log any failures as warnings, and use a finally block to clear the resource reference. This ensures that a failure in closing one resource does not prevent other resources from being closed.
if owns_bqstorage_client and bqstorage_client is not None:
try:
bqstorage_client.close()
except Exception as e:
import warnings
warnings.warn(f"Failed to close bqstorage_client: {e}")
finally:
bqstorage_client = NoneReferences
- When closing multiple resources during cleanup, wrap each close call in a try-except block to handle exceptions individually, log any failures as warnings, and use a finally block to clear the resource references. This ensures that a failure in closing one resource does not prevent other resources from being closed.
Resolve mypy type checking errors in RowIterator for destination table read sessions and Arrow serialization options compatibility across published and editable google-cloud-bigquery-storage stubs. Fix ruff import sorting in bigquery-storage system tests and format branch files with black.
Thank you for opening a Pull Request! Before submitting your PR, there are a few things you can do to make sure it goes smoothly:
Fixes #<issue_number_goes_here> 🦕