samples(bigquery-storage): add Arrow query results samples for query_and_wait and read_rows - #18126
samples(bigquery-storage): add Arrow query results samples for query_and_wait and read_rows#18126alextolpin wants to merge 2 commits into
Conversation
|
Here is the summary of changes. You are about to add 2 region tags.
This comment is generated by snippet-bot.
|
There was a problem hiding this comment.
Code Review
This pull request introduces two new code snippets and their corresponding tests to demonstrate reading BigQuery query results in Apache Arrow format. Specifically, query_and_wait_arrow.py uses the query_and_wait method to fetch Arrow RecordBatches directly, while read_rows_query_job.py initiates a query job and streams the results using the BigQuery Storage Read API. A critical issue was identified in read_rows_query_job.py where the query job is started asynchronously, but the code immediately attempts to read from the stream without waiting for the job to complete. It is recommended to call job.result() to ensure the query finishes before reading.
| # Start the query job. | ||
| job = client.query(query) |
There was a problem hiding this comment.
The client.query(query) method starts an asynchronous query job and returns immediately. Because the job runs asynchronously, attempting to construct the stream name and read from it immediately will fail since the job is still pending or running, and its results are not yet available (additionally, job.location may not be populated yet).
To ensure the query has finished and the results are ready to be read from the stream, you must wait for the job to complete by calling job.result().
| # Start the query job. | |
| job = client.query(query) | |
| # Start the query job and wait for it to complete. | |
| job = client.query(query) | |
| job.result() |
Description
Adds documentation code snippets and system tests demonstrating high-performance query result retrieval in Apache Arrow format with LZ4 compression using the BigQuery Storage API:
query_and_wait()with Arrow format & LZ4 frame compression (query_and_wait_arrow.py):client.query_and_wait()withquery_results_format=enums.QueryResultsFormat.ARROWandcompression_codec=enums.QueryResultsCompressionCodec.LZ4_FRAME.pyarrow.RecordBatchviaresults.to_arrow_iterable().[START bigquerystorage_query_and_wait_arrow].Direct
read_rowson query job default stream (read_rows_query_job.py):BigQueryReadClient.read_rowsagainst the job streamprojects/{project}/locations/{location}/jobs/{job_id}/streams/_default.pyarrow.ipc.[START bigquerystorage_read_rows_query_job].Tests & Dependencies:
query_and_wait_arrow_test.pyandread_rows_query_job_test.pyverifying batch iteration and schema types.pyarrowdependency pins tosamples/snippets/requirements.txt.Follow-up to #18027
Related to #18047
Checklist