Skip to content

[Prototype] Arrow support - #18614

Draft
ohmayr wants to merge 20 commits into
mainfrom
prototype/arrow-support
Draft

ohmayr wants to merge 20 commits into
mainfrom
prototype/arrow-support

Conversation

@ohmayr

@ohmayr ohmayr commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

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:

  • Make sure to open an issue as a bug/issue before writing your code! That way we can discuss the change, evaluate designs, and agree on the general idea
  • Ensure the tests and linter pass
  • Code coverage does not decrease (if any source code was changed)
  • Appropriate docs were updated (if necessary)

Fixes #<issue_number_goes_here> 🦕

ohmayr added 19 commits October 3, 2026 01:03
…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.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines 2649 to +2656
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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

critical

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.

Suggested change
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)

Comment on lines +187 to +201
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,
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

security-medium medium

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.

Suggested change
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
  1. 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.

Comment on lines +2532 to +2533
if owns_bqstorage_client and bqstorage_client is not None:
bqstorage_client._transport.close()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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 = None
References
  1. 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.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant