Reading tables
The library reads tables through the BigQuery Storage Read API. A read runs no query and no job: BigQuery streams the table’s columns as Arrow, and the library decodes them into your structures with serde. When to read a table and when to query it, see Table reads or queries.
use bigquery::*;
use futures::TryStreamExt;
use futures::stream::BoxStream;
use serde::Deserialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const PEOPLE: BigQueryTableId = BigQueryTableId::from_static("people");
#[derive(Debug, Deserialize)]
struct Person {
name: String,
city: String,
year: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let people: BoxStream<BigQueryResult<Person>> = db
.fluent()
.select()
.fields(paths!(Person::{name, city, year})) // Optionally select the columns needed
.from(SHOP.table(PEOPLE))
.filter(|filter| {
filter.for_all([
filter.field(path!(Person::city)).eq("Stockholm"),
filter.field(path!(Person::year)).ge(2010),
])
})
.obj() // Reading rows as structures using serde
.stream_query_with_errors()
.await?;
let as_vec: Vec<Person> = people.try_collect().await?;
println!("{as_vec:?}");
Ok(())
}
A read needs the bigquery.readsessions.create permission on the project and read access to the
table, the same as any Storage Read client.
Selecting columns
.fields(..) takes the column names to read. path! and paths! build them from the fields of
your structure, so a renamed field is a compile error instead of a read that fails on a column the
table does not have:
path!(Person::city)is"city";paths!(Person::{name, city})isvec!["name", "city"];path!(Person::home.county)is"home.county", a subfield of a STRUCT column.
Without .fields(..), a typed read selects by itself the columns that the top-level fields of
your structure name, so it does not read columns it would throw away. It costs one GetTable call
to learn the table’s columns. A structure with #[serde(flatten)], a map, serde_json::Value or
a tuple has no complete list of fields, so it reads every column. record_batches() without
.fields(..) reads every column too.
The paths are Rust field names. A field under #[serde(rename = "...")] needs
path_camel_case! for camelCase columns, or the column name as a plain string.
Be aware that BigQuery checks selected_fields against a schema that lags behind: a column
added or renamed in the last 30 seconds or so is refused, and the read fails with
SchemaMismatchError. The automatic projection falls back to reading every column in that case
and logs a warning.
Filtering
.filter(..) takes a closure that receives a BigQueryFilterBuilder and returns an
Option<BigQueryFilter>, the same shape as in firestore-rs:
f.field(..)witheq,neq,lt,le,gt,ge,is_null,is_not_null,is_inandis_not_in;f.for_allfor AND conditions;f.for_anyfor OR conditions;f.notfor NOT.
You can nest them, and a None entry is dropped, so optional conditions can be written inline:
use bigquery::*;
use serde::Deserialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const PEOPLE: BigQueryTableId = BigQueryTableId::from_static("people");
#[derive(Deserialize)]
struct Person {
name: String,
city: String,
year: i64,
}
async fn example(db: BigQueryDb, min_year: Option<i64>) -> BigQueryResult<()> {
let people: Vec<Person> = db
.fluent()
.select()
.from(SHOP.table(PEOPLE))
.filter(|filter| {
filter.for_all([
filter.for_any([
filter.field(path!(Person::city)).is_in(["Malmö", "Lund"]),
filter.field(path!(Person::city)).is_null(),
]),
min_year.and_then(|year| filter.field(path!(Person::year)).ge(year)),
])
})
.obj()
.query()
.await?;
let _ = people;
Ok(())
}
Storage Read has no query parameters: the filter is SQL text in the session’s row_restriction.
The builder writes values into it as escaped GoogleSQL literals and column names as quoted
identifiers, so a value cannot change the condition it is in, whatever it holds. Values are any
Serialize, with the same mapping as query parameters: a string is a STRING, an integer an
INT64, BigQueryTimestamp a TIMESTAMP, etc. A condition can name any column, selected or not.
A few things to know:
- a NULL value is refused, since
n = NULLmatches no row; useis_null()instead; - GoogleSQL’s three-valued logic applies: a row where the column is NULL matches neither
neq(..)nornot(eq(..)); - a closure that returns
Nonereads every row; - a value without a literal form fails the read before any request, with
SerializeError; - the whole restriction is at most 1 MB, a limit BigQuery checks.
Full example available here.
Raw SQL filters
.filter_sql(..) sends condition text as it is, for example "year > 2010 AND city != ''".
Be aware not to build this text from user input. It is SQL, and a value spliced into it can
rewrite the condition. Use .filter(..) for every condition that carries a value, and keep
.filter_sql(..) for fixed text from your own code. .filter(..) and .filter_sql(..) set the
same restriction, so the last call wins.
Snapshots and sampling
.snapshot_time(..) reads the table as it was at that time instead of now, using BigQuery’s time
travel. It works only within the time travel window of the dataset, 7 days by default:
use bigquery::*;
use serde::Deserialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const PEOPLE: BigQueryTableId = BigQueryTableId::from_static("people");
#[derive(Deserialize)]
struct Person {
name: String,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let an_hour_ago = BigQueryInstant::now() - jiff::SignedDuration::from_hours(1);
let people: Vec<Person> = db
.fluent()
.select()
.from(SHOP.table(PEOPLE))
.snapshot_time(an_hour_ago)
.obj()
.query()
.await?;
let _ = people;
Ok(())
}
.sample_percentage(..) reads a random sample of about that percentage of the table, above 0 and
up to 100. BigQuery samples by storage blocks, so treat the percentage as approximate, especially
on small tables.
Reading rows
.obj::<T>() decodes the rows into T with the library’s own Arrow decoder, which covers every
BigQuery type. Then pick one of:
query()to collect every row into aVec<T>, failing on the first error;stream_query_with_errors()to stream the rows, with every failure as anErritem. A row that fails to decode is oneErr(DeserializeError)and the stream goes on;stream_query()to stream only the rows that decode. Failures are logged aterror!and skipped.
Be aware that a stream from stream_query() that ends does not mean every row was read: a read
stream that failed for good ends it too, only with a log line. Use stream_query_with_errors()
when you need to know.
Rows are decoded on the task that reads their stream, so T is Send + 'static and cannot borrow
from the batch.
Record batches
.record_batches() streams the Arrow RecordBatches as BigQuery sent them, after the IPC decode
and decompression. It is a bit faster than typed rows and the way to hand the data to other Arrow
tools. arrow_array and arrow_schema are re-exported, so you use the same Arrow version as the
library.
BigQueryBatchRows decodes a batch into typed rows later, with the same mapping as .obj().
A row that fails is one Err item, and the rows after it still decode:
use bigquery::*;
use futures::TryStreamExt;
use serde::Deserialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const PEOPLE: BigQueryTableId = BigQueryTableId::from_static("people");
#[derive(Deserialize)]
struct Person {
name: String,
year: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let mut batches = db
.fluent()
.select()
.fields(paths!(Person::{name, year}))
.from(SHOP.table(PEOPLE))
.record_batches()
.await?;
while let Some(batch) = batches.try_next().await? {
println!("{} rows, {} columns", batch.num_rows(), batch.num_columns());
for person in BigQueryBatchRows::<Person>::new(&batch) {
let person = person?;
println!("{} {}", person.name, person.year);
}
}
Ok(())
}
BigQueryBatchRows is not Send, so decode a batch where you hold it and do not keep the iterator
across an .await.
Full example available here.
Parallel streams and resume
One read opens one read session, and BigQuery splits the table into several read streams. Each
stream runs on its own task: ReadRows, the Arrow decode and, for typed reads, the row decode.
The streams meet in one bounded channel, so a slow consumer slows the read down instead of
buffering the table in memory, and rows from different streams arrive in no particular order.
BigQueryReadOptions sets how the session opens:
max_stream_count: the most streams to ask for, the machine’s available parallelism by default;preferred_min_stream_count: the fewest streams BigQuery should aim for;compression:Lz4by default,ZstdorNone.
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const PEOPLE: BigQueryTableId = BigQueryTableId::from_static("people");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let batches = db
.fluent()
.select()
.from(SHOP.table(PEOPLE))
.options(
BigQueryReadOptions::new()
.with_max_stream_count(8)
.with_compression(BigQueryReadCompression::Zstd),
)
.record_batches()
.await?;
let _ = batches;
Ok(())
}
BigQuery decides the real count and often gives fewer streams than asked for. In the benchmarks it gave 4 streams of the 16 asked for on a 1M-row table, and 1 for a 200,000-row query result.
A stream that fails with a retryable error is resumed at its row offset after a backoff, so the
rows it already sent are not read again. Consecutive failures are capped by the client’s max_retries. A
stream that fails for good is one Err item on the read, and then the whole read ends, since a
scan that lost a stream is incomplete.
One case is never retried: an INTERVAL whose time part does not fit Arrow’s nanoseconds fails the
stream on BigQuery’s side, and every resume would fail the same way. Read such a column through a
query with CAST(.. AS STRING) instead.
The read session, its stream count, BigQuery’s estimate of the bytes scanned and the rows and bytes
read are recorded on the BigQuery Read tracing span.