BigQuery for Rust
Library provides a simple API for Google BigQuery using gRPC for every call:
- Fluent high-level and strongly typed API;
- Read and write rows as Rust structures using Serde, with own codecs for every BigQuery type,
including NUMERIC/BIGNUMERIC, INTERVAL, RANGE, nested STRUCT and ARRAY, and JSON columns
mapped by the field type (
String,serde_json::Valueor your own structure); - Support for:
- Table reads through the Storage Read API, in parallel streams, with column projection and typed filters;
- Queries with named and positional parameters, including ARRAY and STRUCT parameters, always bound on the server side;
- Short query mode, so BigQuery can answer short queries without creating a job;
- Large query results read through the Storage Read API automatically;
- Raw Arrow record batches for table reads and query results;
- DML statements with affected row counts, and dry runs;
- Writes through the Storage Write API with built-in batching and backpressure: at least once, exactly once, or atomic (all rows or none);
- Change data capture (CDC): upserts and deletes by the primary key;
- Declarative table schemas, planned and synced with one call: new columns, renames, drops, widening, partitioning, clustering, primary key, recreating an empty table;
- Datasets, tables and jobs management;
- Bytes processed and billed, slot milliseconds and cache hits for every query, as span fields and as results;
- Retries with jitter for the errors where retrying makes sense, including BigQuery rate limits;
- Full async based on Tokio runtime;
- Macros that help you use your structure fields as column paths;
- Google client based on gcloud-sdk library that automatically detects GCE environment or application default accounts for local development;
To start using the library, see Getting started.
Why this crate
Google has its own Rust crate for BigQuery, google-cloud-bigquery (0.18.0 at the time of writing, October 2026). Both crates work with the same APIs in different ways, so below is what each of them provides that the other does not.
What this crate provides
- High-level typed API. Fluent builders in the same style as
firestore-rs. Rows are your own structures in both
directions, through Serde, for table reads, query results, writes and CDC. Dataset and table
IDs are checked types, column paths come from your structure fields with
path!/paths!, and table schemas are declared once and planned or synced with.plan()/.sync(). - gRPC throughout. Queries, jobs, datasets and tables go through the BigQuery v2 API over gRPC, as well as Storage Read and Write. The official crate sends queries over REST and reads their rows as JSON pages.
- Performance. Measured on 2026-10-04 against the official crate and Python on one home
connection in Sweden to
europe-north2, full details in the benchmarks:- small queries take 0.096 s for
SELECT 1, the same as the official crate, which also uses BigQuery’s short query mode, and faster than Python’s 0.173 s, since Python creates a job; - a 200,000-row query result as typed rows takes 1.29 s against 3.86 s, because the library reads a large result through Storage Read and the official crate pages it as JSON;
- a 1M-row table scan as typed rows takes 9.8 s. The official crate has no typed Storage
Read, so its typed path is a
SELECT *query, which took 94.7 s and billed the whole table every run; - a 1M-row Storage Write takes 30.0 s against 32.7 s, with 18% fewer bytes sent as protobuf than the official crate’s Arrow. I think the smaller requests are why it is a bit faster.
- small queries take 0.096 s for
- Observability. Every query, read and write span carries what BigQuery reports: bytes
processed and billed, slot milliseconds, cache hits, rows and bytes read, rows appended and
bytes sent, retries.
query_with_stats()returns a query’s figures together with its rows. - Safety. Query parameter values are always bound on the server side and never written into
the SQL text. The typed filter and the generated DDL write values as escaped literals and names
as quoted identifiers, tested against a corpus of hostile values and on BigQuery itself. The
calls that lose data say so in their names:
dangerously_delete_with_contents(),dangerously_recreate_with_data_loss(). No unsafe code.
What the official crate provides and this one does not
Checked against google-cloud-bigquery 0.18.0 and google-cloud-bigquery-v2 1.0.0:
- It is maintained by Google as part of google-cloud-rust, and its v2 API crate is already 1.0;
- It uses the REST API for v2, which is GA. The v2 API over gRPC this crate uses works for every call the library makes, but Google does not document it and it is pre-GA, so be aware it can change without notice;
- It has a documented client for every v2 service, including models, routines, row access
policies and projects, and every field of the query request (sessions, external tables,
encryption, slot limits, etc.). This crate has its own API for datasets, tables, jobs and the
common query settings, and for the rest only the raw gRPC clients from gcloud-sdk, such as
model_client()androutine_client(); - Its writer takes Arrow record batches and supports buffered streams. This crate writes your structures as protobuf, through the default, committed and pending streams;
- It has stub traits to mock its clients in your tests. This crate has no public mocks.
Crypto provider error
Depends on your other dependencies you may see the error like:
no process-level CryptoProvider available -- call CryptoProvider::install_default() before this point
The TLS crypto providers are not installed by default, so you can choose one. The easiest way to fix it is to include one, for example:
[dependencies]
rustls = "0.23"
If you have several, you may need to call CryptoProvider::install_default() before creating the
client:
rustls::crypto::ring::default_provider().install_default().expect("Failed to install rustls crypto provider");
Getting started
Cargo.toml:
[dependencies]
bigquery = "0.1"
The default feature tls-roots uses the native TLS roots of your system. Use
tls-webpki-roots instead if you want the bundled Mozilla roots:
[dependencies]
bigquery = { version = "0.1", default-features = false, features = ["tls-webpki-roots"] }
Creating a client
BigQueryDb is the client. It opens two authenticated gRPC channels, one to the BigQuery v2 API
for queries, jobs, datasets and tables, and one to the Storage API for reads and writes. Clones
are cheap and share both channels, so create it once and clone it where you need it.
use bigquery::*;
async fn example() -> BigQueryResult<()> {
// Application default credentials, the project given explicitly
let db = BigQueryDb::new("my-gcp-project-id").await?;
// The project detected from GCP_PROJECT, PROJECT_ID or GCP_PROJECT_ID,
// the quota project of the local credentials or the metadata server
let db = BigQueryDb::for_default_project_id().await?;
// With options
let db = BigQueryDb::with_options(
BigQueryDbOptions::new("my-gcp-project-id".to_string()).with_max_retries(5),
)
.await?;
// Everything else starts from the fluent API
let outcome = db.fluent().query("SELECT 1 AS x").execute().await?;
let _ = outcome;
Ok(())
}
for_default_project_id() fails with InvalidParametersError if no project can be detected.
The project ID is checked only for what would break a resource path: it must not be empty and
must not contain / or a control character. Everything else is up to BigQuery.
Client options
BigQueryDbOptions has:
google_project_id: the project that runs the jobs and owns the datasets by default;location: the location to send with jobs and queries, unset by default, see Locations;max_retries: how many times a failed retryable request is sent again,3by default;bigquery_api_url: overrides the v2 API endpoint,https://bigquery.googleapis.com;bigquery_storage_api_url: overrides the Storage API endpoint,https://bigquerystorage.googleapis.com.
Each has a with_... builder method, as in the example above. db.options() returns the options
a client was created with.
Google authentication
Looks for credentials in the following places, preferring the first location found:
- A JSON file whose path is specified by the
GOOGLE_APPLICATION_CREDENTIALSenvironment variable; - A JSON file in a location known to the gcloud command-line tool using
gcloud auth application-default login; - On Google Compute Engine, it fetches credentials from the metadata server.
For local development don’t confuse gcloud auth login with gcloud auth application-default login,
since the first one authorizes only the gcloud tool to access the Cloud Platform.
To use a service account key file directly:
use bigquery::*;
async fn example() -> BigQueryResult<()> {
let db = BigQueryDb::with_options_service_account_key_file(
BigQueryDbOptions::new("my-gcp-project-id".to_string()),
"/path/to/service-account.json".into(),
)
.await?;
let _ = db;
Ok(())
}
For full control over the OAuth2 scopes and the token source there is
with_options_token_source. Its TokenSourceType comes from
gcloud-sdk, which the library does not re-export, so
add gcloud-sdk to your dependencies to use it:
use bigquery::*;
async fn example(service_account_json: String) -> BigQueryResult<()> {
let db = BigQueryDb::with_options_token_source(
BigQueryDbOptions::new("my-gcp-project-id".to_string()),
gcloud_sdk::GCP_DEFAULT_SCOPES.clone(),
gcloud_sdk::TokenSourceType::Json(service_account_json),
)
.await?;
let _ = db;
Ok(())
}
Both channels share one token, so the token source is asked for a new one only when it expires.
The token never reaches the library’s logs or spans. The one info! line at client creation
logs the project, the two endpoints and the scope names.
Endpoints
The library uses Google’s two global endpoints by default. You can change them, for example for a Private Service Connect endpoint or a proxy:
use bigquery::*;
async fn example() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let options = BigQueryDbOptions::new("my-gcp-project-id".to_string())
.with_bigquery_api_url("https://bigquery-myendpoint.p.googleapis.com".parse()?)
.with_bigquery_storage_api_url("https://bigquerystorage-myendpoint.p.googleapis.com".parse()?);
let db = BigQueryDb::with_options(options).await?;
let _ = db;
Ok(())
}
The URLs are url::Url (re-exported as bigquery::url). Creating the client fails with
InvalidParametersError for a scheme other than http or https, or a URL without a host.
There is no emulator support. The BigQuery emulators speak only the REST API, and this library uses gRPC for everything.
The v2 API over gRPC is pre-GA and not documented by Google. It works for every call the library makes, and the library’s live tests run against it, but be aware Google can change it without notice.
Locations
Leave the location unset unless you need it. BigQuery finds the location of a dataset from the
dataset itself and the location of a query from the tables it reads. A wrong location fails with
DataNotFoundError, so setting one only adds a way to fail.
You need it when BigQuery cannot find it by itself, for example for a query that reads no table and should run in a particular region. You can set it for every query and job of a client, or for one query:
use bigquery::*;
const STOCKHOLM: BigQueryLocation = BigQueryLocation::from_static("europe-north2");
async fn example() -> BigQueryResult<()> {
// For every query and job of this client
let db = BigQueryDb::with_options(
BigQueryDbOptions::new("my-gcp-project-id".to_string()).with_location(STOCKHOLM),
)
.await?;
// For one query, overriding the client's location
let outcome = db
.fluent()
.query("SELECT 1")
.location(BigQueryLocation::new("EU")?)
.execute()
.await?;
let _ = outcome;
Ok(())
}
BigQueryLocation is just a name, since Google adds regions all the time. It accepts
regions such as europe-north2 and multi-regions such as US and EU, and checks only that the
name is not empty. Whether BigQuery knows it is checked by BigQuery. from_static checks it at
compile time in a const.
The location BigQuery reports for a job is kept in its BigQueryJobRef, and the job calls
(get_job, cancel_job, delete_job) send it back, so they find jobs in any region without any
setting.
Retries
Requests that fail with a retryable error are sent again, up to max_retries times. The
retryable errors are UNAVAILABLE, RESOURCE_EXHAUSTED, ABORTED, INTERNAL, BigQuery’s rate
limit errors and transport errors. Before retry number n the client waits a random delay of up to
2^(n-1) seconds. Each retry is logged at warn! inside the call’s span.
A retried query does not run a DML statement twice, since every attempt carries the same
request_id. The writer resends batches with its own rules, see Writing data.
Raw gRPC clients
For the calls the library has no API for, BigQueryDb gives you the raw gRPC clients on its
shared channels: job_client(), dataset_client(), table_client(), model_client(),
routine_client(), row_access_policy_client(), project_client(), read_client() and
write_client().
use bigquery::*;
async fn example(db: BigQueryDb) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
use gcloud_sdk::google::cloud::bigquery::v2::GetServiceAccountRequest;
let response = db
.project_client()
.get_service_account(GetServiceAccountRequest {
project_id: db.options().google_project_id.clone(),
})
.await?;
println!("{}", response.into_inner().email);
Ok(())
}
Their types come from gcloud-sdk, which the library does not re-export, so they can change in a patch release.
Running the examples
All examples available in the examples directory. Each one creates its own scratch dataset and deletes it at the end.
To run an example:
PROJECT_ID=<your-google-project-id> cargo run --example query
Queries
The library runs GoogleSQL queries through the v2 Query call over gRPC and reads the results as
your own structures with serde, or as Arrow record batches. When to query and when to read the
table instead, see Table reads or queries.
use bigquery::*;
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct WordCount {
word: String,
word_count: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let words: Vec<WordCount> = db
.fluent()
.query(
"SELECT word, word_count FROM `bigquery-public-data.samples.shakespeare` \
WHERE corpus = @corpus AND word_count >= @min_count \
ORDER BY word_count DESC LIMIT 10",
)
.param("corpus", "hamlet")
.param("min_count", 100)
.obj::<WordCount>()
.query()
.await?;
println!("{words:?}");
Ok(())
}
Columns map to the fields of your structure by name, with the types described in
Type mapping. The target type is DeserializeOwned + Send + 'static, since the rows of a
large result are decoded on each read stream’s task.
Reading the results
.obj::<T>() has these terminals:
query(): every row in aVec, failing on the first error;stream_query_with_errors(): a stream ofBigQueryResult<T>. A row that fails to decode is oneErritem and the stream goes on; a read stream that fails for good is oneErrand then the stream ends;stream_query(): a stream ofT. Failures are logged aterror!and skipped, so a stream that ended does not mean every row was read. Usestream_query_with_errors()if you need to tell the two apart;query_with_stats()andstream_query_with_stats(): the same asquery()andstream_query_with_errors(), with what the job used, see Job stats.
use bigquery::*;
use futures::TryStreamExt;
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct WordCount {
word: String,
word_count: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let mut words = db
.fluent()
.query("SELECT word, word_count FROM `bigquery-public-data.samples.shakespeare`")
.obj::<WordCount>()
.stream_query_with_errors()
.await?;
while let Some(word) = words.try_next().await? {
println!("{}: {}", word.word, word.word_count);
}
Ok(())
}
The query terminal waits for the job to finish before the first row streams. Dropping the stream stops the reading, and does not cancel the job, see Cancellation.
Rows of a large result come from several read streams at once, so they arrive in no particular
order, even with ORDER BY. A result that comes inline keeps its order. If the order matters for a
large result, sort the rows on your side, or ask for one read stream with
.read_options(BigQueryReadOptions::new().with_max_stream_count(1)). I think one stream keeps the
order of the result, Google’s Python client does the same for ORDER BY queries, but the library
does not test it.
To get the result as Arrow, without serde, use record_batches() directly on the query:
use bigquery::*;
use futures::TryStreamExt;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let mut batches = db
.fluent()
.query("SELECT corpus, COUNT(*) AS words FROM `bigquery-public-data.samples.shakespeare` GROUP BY corpus")
.record_batches()
.await?;
while let Some(batch) = batches.try_next().await? {
println!("{} rows, schema {:?}", batch.num_rows(), batch.schema());
}
Ok(())
}
RecordBatch is arrow_array::RecordBatch, re-exported as bigquery::arrow_array.
Parameters
Values go into the query only as query parameters, which BigQuery binds on its side. The library
never writes a value into the SQL text, so a value with quotes, backticks, comments or @x in it
is just a value. A corpus of such values is tested on all the parameter forms below, as STRING,
ARRAY, STRUCT and JSON, against a fake server. Against BigQuery itself a few injection payloads
are tested as STRING parameters.
Be aware this protects only the values: SQL you build with format! from untrusted input is still
your SQL.
Named parameters
.param(name, value) adds @name, with the type inferred from the value’s serde form:
- integers are INT64, floats FLOAT64,
boolBOOL; - strings and
charare STRING, bytes (serde_bytes) BYTES, unit enum variants STRING; - sequences are ARRAY of the first element’s type;
- structures and string-keyed maps are STRUCT, with the fields in order;
- the library’s wrappers are their own types:
BigQueryTimestamp,BigQueryDate,BigQueryTime,BigQueryDateTime,BigQueryJson,BigQueryInterval,BigQueryRange, andBigQueryDecimalas NUMERIC, or BIGNUMERIC when the value does not fit NUMERIC.
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Serialize)]
struct Window {
earliest: BigQueryTimestamp,
latest: BigQueryTimestamp,
}
#[derive(Debug, Deserialize)]
struct Order {
id: i64,
city: String,
}
async fn example(db: BigQueryDb) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let window = Window {
earliest: BigQueryTimestamp("2026-10-01T00:00:00Z".parse()?),
latest: BigQueryTimestamp("2026-10-02T00:00:00Z".parse()?),
};
let orders: Vec<Order> = db
.fluent()
.query(
"SELECT id, city FROM shop.orders \
WHERE city IN UNNEST(@cities) AND placed_at BETWEEN @window.earliest AND @window.latest \
AND total >= @min_total",
)
.param("cities", vec!["Malmö", "Lund"])
.param("window", window)
.param("min_total", BigQueryDecimal("99.90"))
.obj::<Order>()
.query()
.await?;
let _ = orders;
Ok(())
}
Some values have no type to infer: None, an empty sequence, elements of different types. And a
plain jiff value serializes as text, so .param sends it as STRING. For these use
.param_as(name, type, value), which takes the type and accepts the value in any form the
library writes for that type. None is a NULL of that type.
use bigquery::*;
async fn example(db: BigQueryDb, since: Option<jiff::Timestamp>) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query(
"SELECT COUNT(*) FROM shop.orders \
WHERE (@since IS NULL OR placed_at >= @since) \
AND id NOT IN UNNEST(@excluded_ids) \
AND ST_DWITHIN(location, @store, 10000)",
)
.param_as("since", BigQueryFieldType::Timestamp, since)
.param_as("excluded_ids", BigQueryParamType::array_of(BigQueryFieldType::Int64), Vec::<i64>::new())
.param_as("store", BigQueryFieldType::Geography, "POINT(13.0 55.6)")
.execute()
.await?;
let _ = outcome;
Ok(())
}
.params(&value) adds every top-level field of a structure or string-keyed map as a named
parameter, inferred as by .param:
use bigquery::*;
use serde::Serialize;
#[derive(Serialize)]
struct Filter {
city: String,
min_total: f64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let filter = Filter {
city: "Malmö".to_string(),
min_total: 100.0,
};
let outcome = db
.fluent()
.query("DELETE FROM shop.orders WHERE city = @city AND total < @min_total")
.params(&filter)
.execute()
.await?;
let _ = outcome;
Ok(())
}
Parameter names must be GoogleSQL identifiers: ASCII letters, digits and _, not starting with a
digit.
Positional parameters
.positional_param(value) and .positional_param_as(type, value) add the next ?, inferred or
typed the same way as the named ones:
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query("UPDATE shop.orders SET status = ? WHERE id = ?")
.positional_param("shipped")
.positional_param(42)
.execute()
.await?;
let _ = outcome;
Ok(())
}
Named and positional parameters cannot be mixed in one query.
Full example available here.
Errors
The builder methods never fail. A parameter that cannot be encoded is kept in the builder, and the
terminal returns its error, InvalidParametersError or SerializeError, without sending
anything.
Job settings
The query builder also has:
.location(..): where the job runs, see Locations;.default_dataset(..): the dataset unqualified table names resolve in, aBigQueryDatasetIdin the client’s project or aBigQueryDatasetReffor another project;.label(key, value)and.labels(..): labels on the job;.maximum_bytes_billed(bytes): the job fails without running if it would bill more;.use_query_cache(false): BigQuery answers from its query cache by default;.timeout(..): how long the firstQuerycall waits for the job, 10 seconds by default. A job still running then is polled until it completes;.job_timeout(..): how long BigQuery lets the job run before it cancels the job itself;.request_id(..): the idempotency key of theQuerycall, see DML;.inline_rows_limit(..)and.read_options(..): how the rows come back, see Where the rows come from.
use bigquery::*;
use serde::Deserialize;
use std::time::Duration;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
#[derive(Debug, Deserialize)]
struct Order {
id: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let orders: Vec<Order> = db
.fluent()
.query("SELECT id FROM orders WHERE status = 'new'")
.default_dataset(SHOP)
.labels([("team", "shop"), ("report", "new-orders")])
.maximum_bytes_billed(1_000_000_000)
.job_timeout(Duration::from_secs(60))
.obj::<Order>()
.query()
.await?;
let _ = orders;
Ok(())
}
Short query mode
By default the library sends every query with BigQuery’s short query mode
(JOB_CREATION_OPTIONAL), as Google’s own Rust crate does. BigQuery then answers a short query
whose result fits in the first response without creating a job. In the
benchmarks creating the job cost about 65 ms.
A query that ran without a job:
- reports a
query_idand nojobin its stats and outcome; - is still listed in the
INFORMATION_SCHEMA.JOBSviews, with itsquery_idas thejob_idand with its labels; - has nothing for
get_jobto read orcancel_jobto cancel.GetJobrefuses itsquery_id, since it is not a job ID.
BigQuery still creates a job for a query that runs long, a result too large for the response, and DML and DDL statements (every one of them got a job in the library’s live tests).
.job_creation_required() makes every query create a job. Use it when the query needs a job
resource: for job history read through the job calls, for a job ID it is sure to get, or to cancel
it by its job.
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query("SELECT 1")
.job_creation_required()
.execute()
.await?;
let job = outcome.job.expect("a required job is always reported");
println!("{} in {:?}", job.job_id, job.location);
Ok(())
}
A retry of a failed Query call always requires a job, whatever the first attempt asked for.
BigQuery recognises a repeated request_id and returns the job of the first attempt only in that
mode, and answers it with AlreadyExists in the optional one.
DML and DDL
execute() runs a statement, waits for it, and returns a BigQueryQueryOutcome without reading
any rows:
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query("UPDATE shop.orders SET status = 'shipped' WHERE id = @id")
.param("id", 42)
.execute()
.await?;
println!(
"{:?} changed {:?} rows: {:?}",
outcome.statement_type, outcome.num_dml_affected_rows, outcome.dml_stats
);
Ok(())
}
BigQueryQueryOutcome is the same BigQueryJobStats the queries return, see
Job stats. For DML it has:
statement_type:Some(BigQueryStatementType::Update),Insert,Delete,Merge, etc.;num_dml_affected_rows: the rows the statement changed;dml_stats: the rows inserted, updated and deleted, asBigQueryDmlStats.
DDL (CREATE TABLE, ALTER TABLE, etc.) goes through execute() the same way. .obj::<T>() on
a DML or DDL statement is just an empty result.
Every terminal call sends a fresh random request_id unless you set one, and every retry of that
call repeats it, so a retried DML statement is not run twice. BigQuery keeps the keys for a limited
time window. Set your own BigQueryRequestId with .request_id(..) if your code may send the same
statement again on its own, for example after a restart:
use bigquery::*;
async fn example(db: BigQueryDb, order_id: i64) -> BigQueryResult<()> {
let request_id = BigQueryRequestId::new(format!("ship-order-{order_id}"))?;
let outcome = db
.fluent()
.query("UPDATE shop.orders SET status = 'shipped' WHERE id = @id")
.param("id", order_id)
.request_id(request_id)
.execute()
.await?;
let _ = outcome;
Ok(())
}
A statement that fails is an error from the terminal: a syntax error, a missing table or
ERROR() in the SQL come back from the Query call itself, and a job that finished with an error
is a JobError with BigQuery’s reason and messages.
Dry run
dry_run() validates the statement and returns the bytes it would process and the schema of its
result, without running it. It never creates a job and bills nothing.
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let estimate = db
.fluent()
.query("SELECT * FROM `bigquery-public-data.samples.shakespeare`")
.dry_run()
.await?;
println!("Would process {:?} bytes", estimate.total_bytes_processed);
if let Some(schema) = estimate.schema {
for field in schema.fields {
println!("{}: {}", field.name, field.field_type);
}
}
Ok(())
}
Full example available here.
Job stats
query_with_stats() returns the rows together with what the query used, as BigQueryJobStats:
job: the job that ran the query,Nonewhen BigQuery ran it without one;query_id: the ID BigQuery gave the query, with or without a job;statement_type: the kind of statement;total_rows: rows in the result;total_bytes_processedandtotal_bytes_billed: bytes processed, and bytes billed after BigQuery’s rounding and minimums;total_slot_ms: slot milliseconds the job used;cache_hit: whether the query cache answered;num_dml_affected_rowsanddml_stats: for DML.
use bigquery::*;
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct Total {
words: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let (rows, stats) = db
.fluent()
.query("SELECT SUM(word_count) AS words FROM `bigquery-public-data.samples.shakespeare`")
.obj::<Total>()
.query_with_stats()
.await?;
println!(
"{rows:?}: {:?} bytes billed, {:?} slot ms, cache hit {:?}",
stats.total_bytes_billed, stats.total_slot_ms, stats.cache_hit
);
Ok(())
}
A figure BigQuery did not report is None, never 0. The library makes no extra call for the
stats: every figure comes from a response the query receives anyway. A result answered in the
first Query response has only that response’s figures. A larger one also has the figures of the
job’s statistics, which the query reads to find its destination table.
stream_query_with_stats() returns the stats together with the stream, before any row. The query
waits for its job to finish before the first row streams, so every job figure is already known
then. What reading the rows cost is on the read’s span, see Observability.
Even a query that reads no table uses slots. On 1,000 generated rows BigQuery reported 0 bytes processed and billed, and 25 slot ms answered inline, 152 slot ms read through Storage Read.
Full example available here.
Cancellation
Dropping a query’s stream or future does not cancel its job, BigQuery runs it to the end. To limit
how long a job can run, set .job_timeout(..), and BigQuery cancels it itself.
To cancel a job from your code use cancel_job with its BigQueryJobRef. It returns once BigQuery
accepted the request, which is before the job stops, and a job that already finished stays
finished.
A query terminal returns only when its job has finished, so to cancel a long query find it while it runs, for example by a label:
use bigquery::*;
use futures::StreamExt;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
// The long query was started elsewhere with .label("report", "monthly-totals")
let mut running = db
.stream_jobs(BigQueryListJobsParams::new().with_states(vec![BigQueryJobState::Running]))
.await?;
while let Some(job) = running.next().await {
if job.labels.get("report") == Some("monthly-totals") {
db.cancel_job(&job.reference).await?;
}
}
Ok(())
}
A query that runs long always has a job, even in short query mode, since BigQuery creates one for
a query that outlives the first Query call.
Where the rows come from
Every query asks BigQuery for an Arrow result. Then:
- a result that comes complete in the first
Queryresponse is decoded from the inline Arrow in that response; - a larger result, or one whose job outlived the first call, is read from the job’s destination table through the Storage Read API, with several streams in parallel;
- a statement without rows, such as DML or DDL, reads nothing.
Both paths use the same Arrow decoder as table reads, so a type maps the same way in a query and
in a table read. The query’s span records which path it took, as /bigquery/route.
BigQuery decides how much of a result goes inline. A small result reaches its last row sooner inline, since a Storage Read session costs a call to open before the first row. A large result is much faster through Storage Read: in the benchmarks 200,000 rows took 1.29 s, against 3.86 s for the official crate and 4.70 s for Python reading the same result as REST pages.
.inline_rows_limit(rows) caps the rows of the first response, so a result with more goes to
Storage Read. .read_options(..) sets how that read opens its session, the same
BigQueryReadOptions as for table reads:
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let batches = db
.fluent()
.query("SELECT * FROM `bigquery-public-data.samples.shakespeare`")
.inline_rows_limit(10_000)
.read_options(
BigQueryReadOptions::new()
.with_max_stream_count(4)
.with_compression(BigQueryReadCompression::Zstd),
)
.record_batches()
.await?;
let _ = batches;
Ok(())
}
By default the read asks for as many streams as the machine has parallelism, with LZ4 compression, and BigQuery decides how many it actually gives.
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.
Table reads or queries
The library reads data in two ways:
- a table read,
select().from(..)with.filter(..), goes through the Storage Read API. It runs no query and no job; - a query,
query(..)with named parameters, runs GoogleSQL through the v2Querycall.
Both decode rows with the same Arrow decoder, so a type maps the same way in each, see Type mapping.
The same rows both ways:
use bigquery::*;
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,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let read: Vec<Person> = db
.fluent()
.select()
.from(SHOP.table(PEOPLE))
.filter(|filter| filter.field(path!(Person::city)).eq("Malmö"))
.obj()
.query()
.await?;
let queried: Vec<Person> = db
.fluent()
.query("SELECT name, city FROM shop.people WHERE city = @city")
.param("city", "Malmö")
.obj::<Person>()
.query()
.await?;
let _ = (read, queried);
Ok(())
}
Comparison
| Table read | Query | |
|---|---|---|
| What it reads | One table: the columns you select, the rows the filter keeps, at a snapshot time or as a sample | Any GoogleSQL: joins, aggregates, ORDER BY, LIMIT, views, computed columns |
| Writes | No | DML and DDL |
| Values in conditions | Escaped literals in the session’s row_restriction, Storage Read has no parameters | Query parameters, the value never becomes SQL text |
| Job | None, so no job stats, dry run or cancellation | Short query mode answers small ones without a job, others get one |
| Columns read | The ones your structure names, by itself | The ones the SELECT names |
| First row of a small result | After opening a read session | Inline in the first response, about 0.1 s in the benchmarks |
| Large results | Several streams in parallel, each resumed at its row offset after a retryable error | Read from the job’s destination table, through the same Storage Read path |
| Billing | Storage Read: bytes read | Query: bytes processed, then nothing for reading the result |
Be aware that Storage Read cannot read views, logical or materialized, and external tables, as the Storage Read API limitations say. Query them instead.
Billing hints
Prices here are the on-demand list prices in USD from BigQuery pricing in October 2026, so check them for your region and contract:
- a table read is billed $1.10 per TiB read, and the first 300 TiB a month for each billing account are free. Reads within the same location are free of data transfer;
- a query is billed $6.25 per TiB processed, with the first 1 TiB a month free. It bills every
column the query touches over the whole table, unless partitioning or clustering prune it,
even with
LIMIT, and at least 10 MB per table referenced; - reading a query result is free: the result is a temporary table, and reads of those are not billed. Neither is a query answered from the query cache;
- a table read bills the columns it reads, so select only the ones you need. The automatic projection from your structure does it for you, see Selecting columns.
So for a scan of a large table a table read is probably much cheaper, a bit more than a sixth of
the query price per byte, and free under the monthly allowance. In the
benchmarks the 215 MB table scans went through the Storage Read free
tier, while each SELECT * scan of the official crate billed the table. I didn’t measure whether a row
filter reduces the bytes billed for a table read, so I would not count on it.
.maximum_bytes_billed(..) makes a query fail without running if it would bill more, and a
dry run tells you the bytes before you run it.
Which one to use
- Rows of one table, filtered by values: a table read. Exports, syncs, full scans, “the orders of this customer”. It needs no job, reads in parallel and resumes by itself;
- Joins, aggregates, ordering,
LIMIT, views or computed columns: a query, with named parameters for every value; - A small lookup, a few rows by key: probably a query. Short query mode returns it in one call, while a table read opens a session first;
- DML, DDL, or when you need job stats, labels on the job or a dry run: a query;
- Arrow for your own processing: either, both have
record_batches().
Be aware not to splice values into SQL text in either path. Use .filter(..) instead of
.filter_sql(..), and .param(..) instead of format! in the query text.
Writing data
The library writes rows through the BigQuery Storage Write API. Rows are your structures, serialized with serde straight into protobuf against the table’s schema, so there is no JSON and no schema to declare on the client.
There are two ways to write:
db.fluent().insert()for rows you already have: it opens a writer, writes every row, finishes and returns a summary;db.create_streaming_writer()for a long-running producer: one writer you keep open and write rows to as they come.
Both batch the rows by themselves, so you never need to batch before writing.
Inserts
use bigquery::*;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
#[derive(Serialize)]
struct Order {
id: i64,
customer: String,
placed_at: jiff::Timestamp,
}
async fn example(db: BigQueryDb, orders: Vec<Order>) -> BigQueryResult<()> {
// One row
db.fluent()
.insert()
.into(SHOP.table(ORDERS))
.object(&orders[0])
.execute()
.await?;
// Many rows
let summary = db
.fluent()
.insert()
.into(SHOP.table(ORDERS))
.objects(&orders)
.execute()
.await?;
println!("{} rows written", summary.rows_written);
Ok(())
}
objects(..) takes anything iterable whose items are Serialize, so a Vec, a slice or an
iterator over rows built on the fly all work. execute() returns the first failed batch’s error,
or a BigQueryWriteSummary with the rows written, the batches and the bytes sent. A row that does
not serialize stops it with SerializeError. The rows before it may or may not be written: only
the requests already sent can land, and the rows still waiting in a batch are dropped. If you need
all rows or none, use .atomic() or check the rows before the insert.
Every execute() opens a write stream first, which is a round trip of about 300 ms. That is fine
for a load of many rows, but for many small writes keep one streaming writer instead.
Full example available here.
Streaming writer
use bigquery::*;
use futures::StreamExt;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
#[derive(Serialize)]
struct Order {
id: i64,
customer: String,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let (mut writer, mut responses) = db
.create_streaming_writer::<Order>(SHOP.table(ORDERS))
.await?;
// Reading the responses is optional, one item per batch
let responses_task = tokio::spawn(async move {
while let Some(response) = responses.next().await {
match response {
Ok(written) => println!("batch {} written", written.batch_index),
Err(err) => eprintln!("batch failed: {err}"),
}
}
});
for id in 0..10_000 {
let order = Order {
id,
customer: format!("customer-{id}"),
};
writer.write(&order).await?;
}
let summary = writer.finish().await?;
println!("{} written, {} failed", summary.rows_written, summary.rows_failed);
let _ = responses_task.await;
Ok(())
}
The writer owns one AppendRows connection, run by a background task. It is Send but not
Clone, so open more writers if you need more connections. write_all(..) writes several rows in
order.
The response stream yields one BigQueryWriteResponse or one error per batch, in batch order, and
ends when the writer finishes. Reading it is optional, finish() reports the failed rows either
way.
Be aware not to just drop the writer instead of calling finish(). That logs a warning, drops the
rows not acknowledged yet and leaves a pending stream uncommitted.
db.create_streaming_writer_with_options(..) takes BigQueryStreamingWriteOptions, the same
options as .options(..) on an insert.
Full example available here.
Built-in batching
Rows are encoded as you write them and collected into a batch, one AppendRows request. A batch
is sent when one of these comes first:
- size: the next row would take it over
max_request_bytes, 19,000,000 bytes by default; - time:
max_batch_delayhas passed since its first row, 100 ms by default, so the rows of a producer that went quiet still go out; - rows: it holds
max_batch_rowsrows, if you set it, no limit by default; flush()orfinish()is called.
max_request_bytes can be lowered, never raised. BigQuery’s limit is 20 MiB per request, and a
request over it is not rejected on its own: it ends the whole connection. So the library counts
every request’s exact size before sending it. A single row too large for a request alone fails
with SerializeError of kind RowTooLarge and is never sent.
use bigquery::*;
use serde::Serialize;
use std::time::Duration;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
#[derive(Serialize)]
struct Order {
id: i64,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let (writer, _responses) = db
.create_streaming_writer_with_options::<Order>(
SHOP.table(ORDERS),
BigQueryStreamingWriteOptions::new()
.with_max_batch_delay(Duration::from_millis(500))
.with_max_batch_rows(10_000),
)
.await?;
let _ = writer;
Ok(())
}
Backpressure
The writer does not wait for one batch to be acknowledged before sending the next, it pipelines them. Two limits bound what is sent and not acknowledged yet:
max_inflight_requests: 8 by default;max_inflight_bytes: 64 MiB by default.
When either is reached, write().await waits until BigQuery acknowledges a batch. So a producer
faster than the network slows down to its speed instead of buffering rows in memory.
Flush and finish
flush()sends the open batch and waits until every batch written so far has an outcome. It fails only when the writer itself failed for good; a failed batch is reported on the response stream and in the summary.finish()flushes, waits for every acknowledgement and closes the writer, then returns theBigQueryWriteSummary. On the default stream that is all; a committed stream is finalized, and a pending one is finalized and committed.
Write modes
The mode decides which write stream the rows go through, and with it the delivery guarantee:
| Mode | Insert | Writer option | Rows visible | Guarantee |
|---|---|---|---|---|
| Default | .objects(..) | BigQueryWriteMode::Default | as soon as each batch is acknowledged | at least once: a batch resent after a reconnect can be stored twice |
| Exactly once | .exactly_once() | BigQueryWriteMode::Committed | as soon as each batch is acknowledged | exactly once: every request has an offset, and BigQuery recognises a resent one |
| Atomic | .atomic() | BigQueryWriteMode::Pending | all together, at the commit | all rows or none |
| CDC | .changes(..) or .upsert() | default stream | after BigQuery applies the changes | upserts and deletes by primary key, see change data capture |
use bigquery::*;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
#[derive(Serialize)]
struct Order {
id: i64,
}
async fn example(db: BigQueryDb, orders: Vec<Order>) -> BigQueryResult<()> {
// Every row exactly once, even when a request is resent
db.fluent()
.insert()
.into(SHOP.table(ORDERS))
.objects(&orders)
.exactly_once()
.execute()
.await?;
// Every row becomes visible at one commit, or none does
let summary = db
.fluent()
.insert()
.into(SHOP.table(ORDERS))
.objects(&orders)
.atomic()
.execute()
.await?;
println!("committed at {:?}", summary.commit_time);
Ok(())
}
.options(..) replaces the options and the mode with them, so call it before .exactly_once()
or .atomic().
The atomic mode is what a “batch write” in the transactional sense is. If any batch fails,
finish() returns WriteStreamError with the code NOT_COMMITTED and the table gets nothing.
Several pending writers on one table can commit together. finalize() instead of finish()
leaves the commit to db.commit_write_streams(..):
use bigquery::*;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
#[derive(Serialize)]
struct Order {
id: i64,
}
async fn example(db: BigQueryDb, first: Vec<Order>, second: Vec<Order>) -> BigQueryResult<()> {
let options = BigQueryStreamingWriteOptions::new().with_mode(BigQueryWriteMode::Pending);
let (mut first_writer, _) = db
.create_streaming_writer_with_options::<Order>(SHOP.table(ORDERS), options.clone())
.await?;
let (mut second_writer, _) = db
.create_streaming_writer_with_options::<Order>(SHOP.table(ORDERS), options)
.await?;
first_writer.write_all(&first).await?;
second_writer.write_all(&second).await?;
let streams = vec![first_writer.finalize().await?, second_writer.finalize().await?];
let commit_time = db.commit_write_streams(streams).await?;
println!("committed at {commit_time}");
Ok(())
}
Errors
- A row that does not fit the table’s schema fails a streaming writer’s
write()withSerializeError, naming the row by its index in write order. Nothing of it is sent, and the writer goes on. insert().execute()stops at the first row that does not serialize and returnsSerializeError. On the default and committed streams the earlier rows may or may not be written: the requests already sent stay written, and the rows still in the open or queued batches are dropped. With the default batch size a short insert usually has all of its rows in one open batch, so nothing is written. A pending stream (.atomic()) commits nothing.- A batch BigQuery rejects is
RowErrorswith every row BigQuery named, by its index in write order, so you can find and resend them. None of the batch’s rows are written, and the other batches are not affected. - A retryable connection failure makes the writer reconnect and resend every batch not
acknowledged yet, up to the client’s
max_retries. In the default mode that is where a row can be stored twice. - Once the writer failed for good, every
write()returns that error.
Schema changes
The writer follows the table’s schema while it runs. A column added to the table is picked up for the next batches, from BigQuery’s own notice or when a row has a field the writer did not know.
Be aware not to drop a column while a writer still sends it. BigQuery keeps accepting the values for about 9 seconds and drops them silently, before it starts rejecting the rows. Stop every writer from sending the column first, then drop it.
A field a row leaves out is stored as NULL by default. missing_value with
BigQueryMissingValue::DefaultValue asks BigQuery to use the column’s default value expression
instead. The library sends it as BigQuery’s default_missing_value_interpretation and has no live
test for it yet.
Caveats
- A large
.atomic()load stays invisible, and holds its rows in a pending stream, until it is committed. It suits a single load with a clear end. For an endless stream use the default or the exactly once mode. - Rows from two writers on one table can interleave, nothing orders rows between writers.
- Only a CDC sequence number orders changes to the same key, see change data capture.
- Every insert pays the round trip to open its write stream; share a streaming writer for many small writes.
The benchmarks have a Storage Write run of 1M rows against the official crate.
Change data capture
BigQuery tables are made for appending rows. Change data capture (CDC) is BigQuery’s way to keep
a table in sync with a changing source instead: each row you write is a change to the row with
the same primary key, an upsert or a delete, and BigQuery applies the changes for you. There is no
MERGE or UPDATE statement to run and no staging table to clean up.
The typical use is mirroring an OLTP database into BigQuery as a stream. You read the changes of a PostgreSQL or MySQL table from its log, and write each one as it comes. The BigQuery table then looks like the source table, with a delay you choose.
How it works
A CDC table needs a primary key. BigQuery’s keys are NOT ENFORCED: BigQuery never checks them,
so it is up to you that the key is unique in the source. A key can have up to 16 columns.
Every row written through CDC carries two pseudo-columns besides the table’s own columns:
_CHANGE_TYPE:UPSERTinserts the row, or replaces the row with the same key;DELETEdeletes the row with the same key, and only its key columns matter;_CHANGE_SEQUENCE_NUMBER, optional: orders the changes to one key.
Without a sequence number BigQuery orders the changes to one key by the time it received them, the latest wins. That is fine for one producer writing in order, but rows can be resent after a reconnect, and two producers can race. With a sequence number the highest one wins, whenever it arrived, so an old change that arrives late or twice does not overwrite a newer one.
A sequence number is up to four sections of hexadecimal digits separated by /, each up to 16
digits, from 0 to FFFFFFFFFFFFFFFF/FFFFFFFFFFFFFFFF/FFFFFFFFFFFFFFFF/FFFFFFFFFFFFFFFF.
BigQuery compares the sections as numbers, from left to right. A PostgreSQL LSN such as
16/B374D848 is already in this form. If a key gets changes with sequence numbers, send one
with every change to it: mixing changes with and without them on one key gives an unpredictable
order.
BigQuery does not rewrite the table on every change. It keeps the recent changes beside the table
and applies them in the background. The table’s max_staleness option says how old the applied
data may be:
- without
max_staleness, a query merges the changes not applied yet at query time, so it always sees the latest data and pays for that merge; - with
max_staleness = INTERVAL 10 MINUTE, BigQuery applies the changes at least once every 10 minutes with background jobs, and a query reads the applied table, which can be up to 10 minutes old. If the background jobs fall behind the interval, queries merge at query time again.
Either way a reader sees the merged result, never the raw changes. The pseudo-columns cannot be queried.
Creating the table
Declare the primary key with schema(), and set max_staleness with a DDL statement:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const CUSTOMERS: BigQueryTableId = BigQueryTableId::from_static("customers");
struct Customer {
id: i64,
name: String,
city: Option<String>,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
db.fluent()
.schema()
.table(SHOP.table(CUSTOMERS))
.columns(|columns| {
columns.fields([
columns.field(path!(Customer::id)).int64().required(),
columns.field(path!(Customer::name)).string(),
columns.field(path!(Customer::city)).string(),
])
})
.primary_key([path!(Customer::id)])
.cluster_by([path!(Customer::id)])
.sync()
.await?;
db.fluent()
.query("ALTER TABLE shop.customers SET OPTIONS (max_staleness = INTERVAL 10 MINUTE)")
.execute()
.await?;
Ok(())
}
Clustering by the key is what Google’s own CDC example does; it is not required.
Writing changes
.changes(..) on an insert takes BigQueryChanges, each with a BigQueryChangeType and an
optional BigQueryChangeSequenceNumber:
use bigquery::*;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const CUSTOMERS: BigQueryTableId = BigQueryTableId::from_static("customers");
#[derive(Serialize)]
struct Customer {
id: i64,
name: String,
city: Option<String>,
}
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let changes = vec![
BigQueryChange {
change_type: BigQueryChangeType::Upsert,
sequence_number: Some(BigQueryChangeSequenceNumber::from(1)),
row: Customer {
id: 1,
name: "Ada".to_string(),
city: Some("Stockholm".to_string()),
},
},
BigQueryChange {
change_type: BigQueryChangeType::Delete,
sequence_number: Some("16/B374D848".parse()?),
row: Customer {
id: 2,
name: String::new(),
city: None,
},
},
];
db.fluent()
.insert()
.into(SHOP.table(CUSTOMERS))
.changes(changes)
.execute()
.await?;
// Or plain rows as upserts, with no sequence numbers
let customer = Customer {
id: 3,
name: "Grace".to_string(),
city: None,
};
db.fluent()
.insert()
.into(SHOP.table(CUSTOMERS))
.object(&customer)
.upsert()
.execute()
.await?;
Ok(())
}
BigQueryChangeSequenceNumber::from(u64) writes the number in hex, as one section. Parsing a
string takes the text as it is and BigQuery checks its form, so a source that has its own
sequence text, like the LSN above, can pass it through.
For a long-running stream of changes, db.create_cdc_writer(..) opens a CDC writer, the same as
the streaming writer with its batching, backpressure and
response stream:
use bigquery::*;
use serde::Serialize;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const CUSTOMERS: BigQueryTableId = BigQueryTableId::from_static("customers");
#[derive(Serialize)]
struct Customer {
id: i64,
name: String,
city: Option<String>,
}
async fn example(db: BigQueryDb, source_changes: Vec<BigQueryChange<Customer>>) -> BigQueryResult<()> {
let (mut writer, _responses) = db
.create_cdc_writer::<Customer>(
SHOP.table(CUSTOMERS),
BigQueryStreamingWriteOptions::new(),
)
.await?;
for change in &source_changes {
writer.write_change(change).await?;
}
let summary = writer.finish().await?;
println!("{} changes written", summary.rows_written);
Ok(())
}
writer.upsert(&row) and writer.delete(&row) write a change without a sequence number.
Full example available here.
Why protobuf
BigQuery takes CDC changes only as protobuf rows on the default write stream, Apache Arrow rows are not supported for CDC. This is the main reason the library writes every row as protobuf: one encoder serves plain inserts and CDC alike. Protobuf requests were also smaller than Arrow for the same rows in the benchmarks, 18% fewer bytes for 1M rows.
Limits
From BigQuery’s side, as its CDC limitations list them:
- CDC goes only through the default stream, so
.exactly_once()and.atomic()cannot be combined with it; the library refuses them withInvalidParametersErrorbefore any request; - the primary key is not enforced and has at most 16 columns;
- the table has at most 2,000 top-level columns;
- mutating DML (
UPDATE,DELETE,MERGE), wildcard table queries and search indexes are not supported on the table; - while queries merge at query time because
max_stalenessis too low or not set, the table cannot be copied, cloned or snapshotted, and cannot be read through the Storage Read API, so table reads need amax_stalenessthe background jobs keep up with; - a query that merges at query time scans the whole table, whatever its partition filter;
- exports do not include changes not applied yet;
- the applying is BigQuery compute and is billed as such, on demand unless you have a BACKGROUND reservation, which Standard edition does not have.
From the library’s side:
- a CDC writer always writes
_CHANGE_TYPE, so write plain inserts to the same table with a separate writer. BigQuery’s CDC documentation says rows with and without a change type on one connection are not supported; - a delete still serializes a whole row of the table’s type, but only its key columns matter;
- delivery is at least once, as on every default stream write. With sequence numbers a change sent twice is harmless.
Schema management
Tables are usually created with DDL in the console or a migration script, and changed by hand when the code needs a new column. The library supports declaring a table’s schema and settings in Rust instead, next to the structs that read and write it, and making the table match with one explicit call.
It works the same way as the index management in
firestore-rs: a declaration, a read-only .plan()
and a .sync() that applies it.
Declaring a table
Everything starts from db.fluent().schema().table(..), then a chain of declarations ending in
.plan() or .sync():
use bigquery::*;
struct Address {
city: String,
}
struct Order {
id: i64,
customer: String,
total: String,
shipping: Address,
tags: Vec<String>,
placed_at: jiff::Timestamp,
}
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let report = db
.fluent()
.schema()
.table(SHOP.table(ORDERS))
.columns(|columns| {
columns.fields([
columns.field(path!(Order::id)).int64().required(),
columns.field(path!(Order::customer)).string_with_max_length(64),
columns.field(path!(Order::total)).numeric_with(10, 2).default_value("0"),
columns.field(path!(Order::shipping)).record(|address| {
address.fields([address.field(path!(Address::city)).string()])
}),
columns.field(path!(Order::tags)).string().repeated(),
columns.field(path!(Order::placed_at))
.timestamp()
.description("When the customer placed the order"),
])
})
.primary_key([path!(Order::id)])
.partition_by_day(path!(Order::placed_at))
.cluster_by([path!(Order::customer)])
.description("Orders")
.labels([("team", "shop")])
.sync()
.await?;
println!("{report}");
Ok(())
}
A missing table is created with one InsertTable. An existing one is changed in place where
BigQuery allows it.
.columns() declares the columns in table order. Each column needs one type:
int64(),float64(),bool();numeric(),numeric_with(precision, scale),bignumeric(),bignumeric_with(..);string(),string_with_max_length(n),bytes(),bytes_with_max_length(n);date(),time(),datetime(),timestamp(),interval(),range(element);geography(),json();record(|address| ..)for a STRUCT, with its fields declared the same way;of_type(BigQueryFieldType)for any of the above as a value.
A column is NULLABLE until required() or repeated() says otherwise. description(..) and
default_value(..) are optional, and renamed_from(..) is described below.
The table settings are primary_key(..) (always NOT ENFORCED, as every BigQuery key),
partition_by_hour/day/month/year(column) or partition_by(BigQueryPartitioning),
partition_expiration(Duration), cluster_by(..), description(..), labels(..) and
expiration(BigQueryInstant).
Be aware a default value is trusted SQL. It is sent to BigQuery as it is, and written into
generated DDL as one parenthesized expression, so it must never carry text from your users. A
string default is written as its own quoted literal, default_value("'none'"). Descriptions,
labels and option values are always sent as values or escaped literals.
The library checks only what it needs itself to build the requests: a column without a type, an
empty column name or one with . or a control character, two columns that differ only in case,
renamed_from on a nested field, a partitioning column that is not declared, cluster_by or
primary_key with no columns, partition_expiration without partitioning and snapshot_first()
without a recreate opt-in. These fail .plan() and .sync() with InvalidParametersError
before any request. Lengths, naming rules and label rules are left to BigQuery.
plan() versus sync()
.plan() is read-only. It sends one GetTable (and ListRowAccessPolicies when a recreate is
planned) and reports what .sync() would do, writing nothing.
.sync() plans the same way and then applies the plan.
Both return a value with a Display impl, so it can be printed or logged directly:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let plan = db
.fluent()
.schema()
.table(SHOP.table(ORDERS))
.columns(|columns| {
columns.fields([
columns.field("id").int64().required(),
columns.field("customer").string(),
columns.field("note").string(),
])
})
.plan()
.await?;
if !plan.is_empty() {
println!("{plan}");
}
Ok(())
}
BigQueryTablePlan has:
create: the table to create, when it does not exist;changes: the in-place changes, in the order.sync()sends them;withheld: changes left out until the declaration opts in to them, each with itsBigQueryWithheldReason(PruneUndeclaredorAllowWidening);recreate: the recreate.sync()runs instead ofchanges, when a change is impossible in place and a recreate opt-in applies;impossibleandrefusal: the changes that are impossible in place, and why.sync()refuses them.
BigQueryTableSyncReport from .sync() has the same shape for what was done: created,
applied, withheld, recreated, snapshot and dropped_data, the last one listing every
column or row a sync deleted.
One statement owns one table
A chain names exactly one table, and that table is the unit of ownership: .prune_undeclared()
never reaches anything outside it. Several tables need several statements, for example one per
table in a startup function.
What a declaration leaves out
What the chain declares is what .sync() makes the table hold. What it leaves undeclared is kept
as the table has it:
- an undeclared column, label, clustering, primary key or partitioning is kept and listed under
withheldwithPruneUndeclared; - an undeclared description, default value or expiration is kept and not reported.
So a declaration can start small, with only the columns your code needs, and the rest of the table stays untouched.
Before comparing, both sides are normalised the way BigQuery itself compares them. GetTable
returns the legacy type names, so INTEGER equals INT64, FLOAT equals FLOAT64, RECORD equals
STRUCT, etc. A column created by DDL comes back with an empty mode, which equals NULLABLE. Type
parameters are part of the type, so STRING(10) and STRING(20) differ, and NUMERIC(10) equals
NUMERIC(10, 0). Column names compare ignoring case.
Change classes
Every difference between the declaration and the table is one BigQuerySchemaChange, and BigQuery
applies each kind in its own way:
| Change | Applied by | Needs |
|---|---|---|
| add a NULLABLE or REPEATED column, a field inside an existing RECORD, a RECORD column | PatchTable | |
| add a column with a default value | PatchTable, then a second PatchTable for the default | |
| relax REQUIRED to NULLABLE | PatchTable | |
| set a column default, a column or table description | PatchTable | |
| add or change a label, the expiration, the partition expiration | PatchTable | |
| add or change clustering, add or change the primary key | PatchTable | |
| remove the primary key | PatchTable | prune_undeclared() |
| remove a label or the clustering | UpdateTable | prune_undeclared() |
| rename a column | DDL RENAME COLUMN | renamed_from(..) |
| widen a column type | DDL ALTER COLUMN SET DATA TYPE | allow_widening() |
| drop a column, with every value in it | DDL DROP COLUMN | prune_undeclared() |
| anything else | impossible in place | a recreate opt-in |
BigQuery cannot add a column with a default in one step, so the column is added first and the default set right after. Rows already in the table stay NULL in that column; only new rows get the default. The same holds for a default set on an existing column.
Impossible in place means BigQuery has no call or statement that does it without recreating the table:
- a type change that is not a widening, INT64 to FLOAT64 and every narrowing included;
- a change between a RECORD and another type, or a widening of a nested field;
- NULLABLE to REQUIRED, and any change to or from REPEATED;
- a new REQUIRED column or nested field;
- dropping a nested field;
- adding or changing the partitioning, or removing it with
prune_undeclared().
Without a recreate opt-in .sync() refuses such a plan and writes nothing, not even the changes
that are possible in place. The error is SchemaChangeRefused, carrying the plan:
use bigquery::*;
use bigquery::errors::BigQueryError;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let synced = db
.fluent()
.schema()
.table(SHOP.table(ORDERS))
.columns(|columns| columns.fields([columns.field("id").string().required()]))
.sync()
.await;
match synced {
Ok(report) => println!("{report}"),
Err(BigQueryError::SchemaChangeRefused(refused)) => {
for change in &refused.plan.impossible {
eprintln!("impossible in place: {change}");
}
}
Err(err) => return Err(err),
}
Ok(())
}
The order sync() writes changes in
.sync() reads the table once and writes in a fixed order:
- one
PatchTablewith every change that can go into a patch; - a second
PatchTablewith the default values of the columns the first one added; - one
UpdateTablefor removed labels and clustering, since a patch cannot remove them; - DDL through a query, one statement at a time: renames, then widenings, then drops.
Patch and Update carry the etag of the table the sync read, as an if-match precondition. When the
table changed in between, the write fails with DataConflictError and nothing after it is sent.
Run the sync again: it reads the table as it is now. Changes sent before a failure stay applied, and
the sync logs at warn what it applied so far.
BigQuery limits how often one table can be updated: about 5 DDL statements or 7 to 8 patches in a quick burst, then it starts rejecting them. So a sync sends its first five writes back to back and spaces the rest about 2 seconds apart, and the rate-limit errors are retried.
Rolling out a change
BigQuery applies a schema change at once, but running writers and readers see it later. During a rolling update old and new replicas run side by side, so a sync that removes a column before the old replicas stop is the same as removing it under them.
So run .sync() once per release, in two steps, the same as firestore-rs’s index management:
- before the deploy, the sync without
.prune_undeclared(): it adds columns, relaxes REQUIRED, sets descriptions, labels, clustering, etc., all of which keep running writers working; - after the rollout has finished, the same sync with
.prune_undeclared(), which drops what the release no longer declares.
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
fn orders(db: &BigQueryDb) -> BigQueryTableSchemaBuilder<'_> {
db.fluent()
.schema()
.table(SHOP.table(ORDERS))
.columns(|columns| {
columns.fields([
columns.field("id").int64().required(),
columns.field("customer").string(),
columns.field("channel").string(),
])
})
}
async fn example(db: BigQueryDb, rollout_finished: bool) -> BigQueryResult<()> {
let report = if rollout_finished {
orders(&db).prune_undeclared().sync().await?
} else {
orders(&db).sync().await?
};
println!("{report}");
Ok(())
}
From a Kubernetes Job or a CI/CD step, run the first one before the deploy step and the second one
once the new version is serving.
What writers and readers see after a sync:
- New columns. New writer connections accept rows with the new column about a second after the
sync. Open connections get BigQuery’s
updated_schemaafter about 7 seconds, and the library’s streaming writer encodes the next batches against it. Nothing is lost on the way. - Relaxed columns. An open connection can keep rejecting NULL for a relaxed column after the sync: for 2.1 s in one test and still after 5 minutes in another. A fresh connection accepted NULL after 0.4 s and 11.5 s. The library’s streaming writer reconnects by itself once it sees the relaxed column, so you only need to reconnect other writers before they send NULL.
- Dropped columns. For about 9 seconds after a drop, an open writer’s values for the column are
accepted and lost silently. Then the connection fails with
Input schema has more fields than BigQuery schema. That is why a drop waits for the second step. - Readers. A read with selected fields that names a new or renamed column fails with
The following selected fields do not exist in the table schemafor about 30 seconds after the sync. A read without selected fields sees the change at once.
Renaming a column
renamed_from(old_name) declares that a column used to be called old_name:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
db.fluent()
.schema()
.table(SHOP.table(ORDERS))
.columns(|columns| {
columns.fields([
columns.field("id").int64().required(),
columns.field("customer_name").string().renamed_from("customer"),
])
})
.sync()
.await?;
Ok(())
}
When the table has customer and no customer_name, .sync() renames it with RENAME COLUMN and
the values are kept. Once the table has customer_name, the declaration is a no-op, so it can stay
in the code for a while.
For writers a rename is a drop of the old name: a writer still sending customer loses that value
silently for a few seconds and then fails. renamed_from(..) is not held back by
.prune_undeclared(), so the rename runs in whichever sync carries it. Be aware not to add
renamed_from(..) to the declaration the first step runs: add it only to the one used after the
rollout, once the writers have moved to the new name. Only top-level columns can be renamed.
Full example available here.
Widening a column
A declared type that is wider than the column’s is withheld with AllowWidening until
.allow_widening() is set. Then .sync() widens it with ALTER COLUMN SET DATA TYPE, and the
values are kept. The library takes these pairs as widenings:
- INT64 to NUMERIC or BIGNUMERIC, and NUMERIC to BIGNUMERIC, without parameters;
NUMERIC(p, s)orBIGNUMERIC(p, s)to parameters that fit every value of the old ones, or to no parameters;STRING(n)andBYTES(n)to a longer maximum length, or to no maximum.
INT64 to NUMERIC and a longer STRING(n) were checked against BigQuery; the other pairs follow
GoogleSQL’s assignability rules. INT64 to FLOAT64 is not a widening for BigQuery, and every other
type change is impossible in place. Only top-level columns can be widened.
Removing what you no longer declare
.prune_undeclared() makes .sync() also remove what the table has and the statement does not
declare:
- undeclared columns are dropped, with every value in them;
- undeclared labels, clustering and the primary key are removed.
Be aware this is the one in-place change that loses data. Every dropped column is logged at warn
and listed in the report’s dropped_data. An undeclared nested field or partitioning cannot be
removed in place, so with .prune_undeclared() set they make the plan need a recreate.
Recreating a table
When a change is impossible in place, the only way is to drop the table and create it again, losing its rows. No data is migrated. A declaration opts in to that with one of:
recreate_if_empty(): recreate only whenGetTablereportsnum_rows == 0, otherwise refuse withBigQueryRefusal::NotEmpty;dangerously_recreate_with_data_loss(): recreate whatever the table holds.
The last of the two in a chain wins. snapshot_first() adds a CREATE SNAPSHOT TABLE before the
recreate, as a cheap undo:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const EVENTS: BigQueryTableId = BigQueryTableId::from_static("events");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let report = db
.fluent()
.schema()
.table(SHOP.table(EVENTS))
.columns(|columns| {
columns.fields([
columns.field("id").string().required(),
columns.field("happened_at").timestamp(),
])
})
.partition_by_month("happened_at")
.dangerously_recreate_with_data_loss()
.snapshot_first()
.sync()
.await?;
if let Some(snapshot) = &report.snapshot {
println!("snapshot taken: {snapshot}");
}
Ok(())
}
The table is recreated as declared, plus whatever it has undeclared and not pruned. There are two ways, and the library picks one by what changes:
CREATE OR REPLACE TABLE | DROP TABLE and CREATE TABLE in one script | |
|---|---|---|
| used when | the partitioning and clustering stay the same | the partitioning or clustering changes, which CREATE OR REPLACE cannot do |
| rows | lost | lost |
| columns, key, partitioning, clustering, description, labels, expirations | as declared or kept | as declared or kept |
| table IAM bindings | kept | lost |
| row access policies | lost | lost |
The plan lists the row access policies that will be lost, and BigQueryRecreate names the method.
The snapshot is named <table>_snapshot_<unix seconds> in the same dataset, and the library never
deletes it.
These are the limits of a recreate, as measured against BigQuery, and the library does not hide any of them:
num_rowslags Storage Write. Rows written through committed or pending streams show innum_rowsafter 60 to 70 seconds, and rows on the default stream can take more than 5 minutes. So withrecreate_if_empty()a table written seconds ago can read as empty and be recreated, its rows lost. The row and byte counts in the plan and indropped_datahave the same lag and can be lower than what is lost. Userecreate_if_empty()for tables that are new or being experimented on, with no writer running.- The snapshot misses the streaming buffer. A
CREATE SNAPSHOT TABLEdoes not hold rows still in the streaming buffer, which is every row written through Storage Write in roughly the last minute or more. - Nothing guards a recreate. BigQuery takes no precondition on
DROP TABLEorCREATE OR REPLACE, so a change made to the table between the plan and the recreate is lost too. - Writers do not see it. A writer connection opened before the recreate keeps getting
successful acks for about 4 to 7 seconds, for rows that are lost. Then the next append waits about
120 seconds for
DeadlineExceeded. New writers can getNotFound(is truncated,is re-created) for up to a few minutes, so the report does not promise that the table is writable yet. A writer still using the old column types getsInvalidArgumentat once, and that one is not retryable. - Readers do not see it either. A read session opened before the recreate keeps returning the old rows and the old schema. Readers have to start a new session.
So stop writers and readers before a recreate. A dangerous recreate is logged at warn with the row
count it read.
Full example available here.
Logs and spans
Every .plan() and .sync() runs in a BigQuery schema span with the table in the
/bigquery/table field. A created table and every applied write are logged at info, and
the data-losing steps at warn: a dropped column, a dangerous recreate and a sync that failed part
way, with the report of what it had applied.
Datasets, tables and jobs
Besides the declarative table schemas, the library provides the plain admin calls over the v2 API:
- create, read, update, delete and list datasets;
- read, delete and list tables;
- read, delete, cancel and list jobs.
Every call returns the library’s own types, such as BigQueryDataset, BigQueryTable and
BigQueryJob, and a value BigQuery leaves out is None.
Dataset and table IDs
Every call that names a dataset or a table takes the validated ID types, BigQueryDatasetId and
BigQueryTableId. Declare the ones your application knows up front as constants, and build
references from them:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
// The table in the client's project
let orders = SHOP.table(ORDERS);
assert_eq!(orders.to_string(), "shop.orders");
// The same table in another project
let other = BigQueryDatasetRef::new("acme-prod", SHOP)?.table(ORDERS);
assert_eq!(other.to_string(), "acme-prod.shop.orders");
Ok::<(), bigquery::errors::BigQueryError>(())
from_static checks the literal at compile time, so an invalid one in a const fails the build.
An ID that arrives at run time goes through new, parse() or try_into(), and deserializing an
ID checks it the same way:
use bigquery::*;
let dataset = BigQueryDatasetId::new("shop_eu")?;
let table: BigQueryTableId = "orders_2026".parse()?;
let orders = dataset.table(table);
let _ = orders;
assert!(BigQueryDatasetId::new("shop.eu").is_err());
Ok::<(), bigquery::errors::BigQueryError>(())
The check is deliberately narrow. An ID is rejected when it is empty, longer than 1,024 bytes, or
holds a control character or any of `, ', ", \, ., /, $ and @, since those
would change a resource path, a dotted reference or SQL text built from the ID. BigQuery’s full
naming rules are left to BigQuery, so an ID that passes here can still be refused by the call that
sends it.
BigQueryTableRef and BigQueryDatasetRef also parse from text, dataset.table or
project.dataset.table. Be aware not to parse text you do not trust this way, since the text picks
the project. Build the reference from the IDs instead, or check that project() is None after
parsing.
The project ID stays a plain String, because it is shared by every Google Cloud product. It is
checked only for being non-empty and free of / and control characters. A location is a
BigQueryLocation, such as BigQueryLocation::from_static("EU"), checked only for being
non-empty.
Datasets
Datasets are reached through db.fluent().schema().dataset(..), which takes a BigQueryDatasetId
for the client’s project or a BigQueryDatasetRef for another one:
use bigquery::*;
use std::time::Duration;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const EU: BigQueryLocation = BigQueryLocation::from_static("EU");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let shop = db
.fluent()
.schema()
.dataset(SHOP)
.create()
.location(EU)
.description("Shop data")
.default_table_expiration(Duration::from_secs(90 * 24 * 60 * 60))
.labels([("team", "shop")])
.execute()
.await?;
println!("created {} in {:?}", shop.reference, shop.location);
let shop = db.fluent().schema().dataset(SHOP).get().await?;
println!("{:?}", shop.labels.get("team"));
Ok(())
}
The location cannot change after the dataset is created. Unset, BigQuery picks US. Creating a
dataset that already exists fails with DataConflictError.
update() changes the description and the labels, and keeps everything else as the dataset has
it:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let shop = db
.fluent()
.schema()
.dataset(SHOP)
.update()
.label("stage", "prod")
.remove_label("tmp")
.clear_description()
.execute()
.await?;
let _ = shop;
Ok(())
}
description(..)sets the description,clear_description()removes it;label(key, value)adds or changes one label,remove_label(key)removes one, andlabels(..)replaces all of them. Label changes apply in the order they are made.
The update reads the dataset and writes it back with the etag it read as an if-match
precondition, so a dataset that changed in between fails with DataConflictError and nothing is
written. Run it again. Only the metadata is written, the access list stays as it is.
There are two ways to delete a dataset:
delete()deletes an empty dataset. BigQuery refuses one that holds any table, view, model or routine, and nothing is deleted;dangerously_delete_with_contents()deletes every table in the dataset with every row in it, and its views, models and routines, then the dataset. Nothing is kept for an undo.
use bigquery::*;
const SCRATCH: BigQueryDatasetId = BigQueryDatasetId::from_static("scratch");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
db.fluent()
.schema()
.dataset(SCRATCH)
.dangerously_delete_with_contents()
.await?;
Ok(())
}
Listing datasets and tables
Listings are streams that fetch the next page when the stream reaches it:
use bigquery::*;
use futures::StreamExt;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let mut datasets = db.fluent().schema().datasets().stream_all().await?;
while let Some(dataset) = datasets.next().await {
println!("{} in {:?}", dataset.reference, dataset.location);
}
let tables: Vec<BigQueryTableSummary> = db
.fluent()
.schema()
.dataset(SHOP)
.tables()
.page_size(100)
.stream_all()
.await?
.collect()
.await;
let _ = tables;
Ok(())
}
stream_all() logs a failed page at error and ends there, since the listing cannot go on without
the token that page would have returned. stream_all_with_errors() yields the failure as the
stream’s last item instead. page_size(..) sets how many items one call returns, and
datasets().project(..) lists another project’s datasets.
The listings return summaries, BigQueryDatasetSummary and BigQueryTableSummary, with what
BigQuery’s list calls return: the reference, the location or table type, the labels, etc. Read the
full value with get().
Full example available here.
Tables
Tables are created and changed through the declarative schemas, see Schema management. The same builder has two plain calls, which act on the table and ignore anything declared before them:
use bigquery::*;
const SHOP: BigQueryDatasetId = BigQueryDatasetId::from_static("shop");
const ORDERS: BigQueryTableId = BigQueryTableId::from_static("orders");
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let orders = db.fluent().schema().table(SHOP.table(ORDERS)).get().await?;
for field in &orders.schema.fields {
println!("{} {} {}", field.name, field.field_type, field.mode);
}
println!("{:?} rows, partitioning {:?}", orders.num_rows, orders.partitioning);
db.fluent().schema().table(SHOP.table(ORDERS)).delete().await?;
Ok(())
}
get() returns a BigQueryTable: the reference, the table_type, the schema as
BigQueryTableSchema, the description, labels, partitioning, clustering, num_rows, num_bytes,
the location, and the creation, last modified and expiration times. It reads views, materialized
views, external tables and snapshots as well. The schema uses the same types as the
type mapping, so INTEGER from the v2 API reads as BigQueryFieldType::Int64.
Be aware num_rows and num_bytes lag Storage Write: rows written through committed or pending
streams show after about a minute, and rows on the default stream can take more than 5 minutes. A
SELECT COUNT(*) gives the current count.
delete() deletes the table with every row in it. Nothing is kept for an undo, and open writers to
it fail.
Jobs
Every query that runs as a job names it in its outcome, as a BigQueryJobRef with the project,
the job ID and the location. The job calls are methods on BigQueryDb:
use bigquery::*;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query("SELECT 1")
.job_creation_required()
.execute()
.await?;
if let Some(job_ref) = &outcome.job {
let job = db.get_job(job_ref).await?;
println!(
"{:?} {:?}, billed {:?} bytes",
job.statement_type, job.state, job.total_bytes_billed
);
db.delete_job(job_ref).await?;
}
Ok(())
}
get_job(&job_ref)reads a job: its type, state, error, user, labels, statement type, the creation, start and end times, and the bytes processed and billed;cancel_job(&job_ref)asks BigQuery to cancel a running job and returns once the request is accepted, before the job stops. Dropping a query’s stream or future does not cancel its job;delete_job(&job_ref)deletes the metadata of a finished job, for example a failed query’s SQL that should not stay in the job history. It does not cancel a running job.
A short query that BigQuery answered without a job names none. .job_creation_required() makes
the query always run as a job.
stream_jobs(..) lists jobs, newest first, with the filters of BigQueryListJobsParams:
use bigquery::*;
use futures::StreamExt;
async fn example(db: BigQueryDb) -> BigQueryResult<()> {
let since: BigQueryInstant = "2026-10-01T00:00:00Z".parse().expect("valid instant");
let mut jobs = db
.stream_jobs(
BigQueryListJobsParams::new()
.with_min_creation_time(since)
.with_states(vec![BigQueryJobState::Done]),
)
.await?;
while let Some(job) = jobs.next().await {
println!("{} {:?}", job.reference.job_id, job.total_bytes_billed);
}
Ok(())
}
The filters are the project, all_users (which needs the Owner role on the project), the minimum
and maximum creation time, the states, the parent job of a script and the page size.
stream_jobs_with_errors(..) yields a failed page as the last item, the same as for the other
listings.
The job state, the job type, the statement type and the table type are enums with an
Other(String) case for names the library does not know yet, and as_str() gives back the name
BigQuery sent for every value.
Write streams
There is no listing of write streams. The Storage Write API has no list call, so a stream is reachable only by the name its creator got back.
Errors and retries
- A missing dataset, table or job is
DataNotFoundError; - creating a dataset that exists, and an update that lost the
if-matchprecondition, areDataConflictError; - every call goes through the client’s retries.
A retry after a lost response repeats a call that may already have succeeded. So a create can
report DataConflictError for the dataset it created, and a delete DataNotFoundError for the
dataset, table or job it deleted.
Type mapping
The library reads and writes rows with serde, through its own codecs:
- reads, from Storage Read and from inline query results, go through one Arrow decoder;
- writes, through Storage Write, go through one protobuf encoder.
Both codecs share one mapping, so every BigQuery type has one set of Rust forms, the same in both directions. Whatever you write in a form, you can read back in the same form.
The library re-exports jiff, so the temporal types need no extra dependency.
Modes
A column’s mode decides the Rust shape around its type:
| NULLABLE | REQUIRED | REPEATED | |
|---|---|---|---|
| Rust field | Option<T> | T | Vec<T> |
| also on read | a bare T, which fails a NULL row with NullForNonOption | Option<T>, always Some | Option<Vec<T>>, always Some |
None on write | stored as NULL | NullForRequired | an empty array |
| field missing from the Rust row on write | stored as NULL | MissingRequiredField | an empty array |
| NULL element on write | NullArrayElement |
BigQuery never stores a NULL array: an empty REPEATED column reads as an empty Vec. A REPEATED
column can be any serde sequence on both sides, such as VecDeque<T>, BTreeSet<T> or [T; N]
when the array has exactly N elements.
An ARRAY column is the REPEATED mode of its element type, so there is no ARRAY type of its own.
BigQuery has no array of arrays as a column type, and the encoder refuses one with
UnsupportedType.
Every type at a glance
“Default” works with plain #[derive(Serialize, Deserialize)] and no attributes. The other forms
work in both directions too.
| BigQuery type | Default Rust type | Other forms | Wrapper |
|---|---|---|---|
| INT64 | i64 | other integer types, range-checked | |
| FLOAT64 | f64 | f32 | |
| NUMERIC, BIGNUMERIC | String | integers, f64 | BigQueryDecimal<T> |
| BOOL | bool | ||
| STRING | String | Box<str>, char, unit enums, any type serialized as a string | |
| BYTES | Vec<u8> | serde_bytes::ByteBuf, [u8; N] | |
| DATE | jiff::civil::Date | String, i32 days | BigQueryDate |
| TIME | jiff::civil::Time | String, i64 microseconds of the day | BigQueryTime |
| DATETIME | jiff::civil::DateTime | String, i64 civil microseconds | BigQueryDateTime |
| TIMESTAMP | jiff::Timestamp | String, i64 microseconds since the epoch | BigQueryTimestamp |
| GEOGRAPHY | String (WKT) | ||
| JSON | by the Rust type, see JSON | BigQueryJson<T> | |
| INTERVAL | BigQueryInterval | String | |
| RANGE | BigQueryRange<T> | ||
| STRUCT | a struct with derived serde | HashMap<String, V>, BTreeMap<String, V>, serde_json::Value |
A form not in this table is not part of the mapping, even if it happens to work. A form the column
does not take fails that row with TypeMismatch, in either direction.
A typical row looks like this:
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
struct Order {
id: i64,
customer: Option<String>,
total: String,
paid: bool,
tags: Vec<String>,
placed_on: jiff::civil::Date,
placed_at: Option<jiff::Timestamp>,
shipping: Option<Address>,
attributes: serde_json::Value,
}
#[derive(Debug, Serialize, Deserialize)]
struct Address {
city: String,
street: Option<String>,
}
Integers and floats
INT64 reads into any Rust integer type, i8 to i128 and u8 to u128, and a value that does not
fit fails that row with OutOfRange. On write a value above i64::MAX is OutOfRange.
FLOAT64 is f64, or f32, which is cast on read with no range check. NaN, both infinities and
-0.0 round-trip exactly. Integers are not taken for FLOAT64 and floats are not taken for INT64,
in either direction.
NUMERIC and BIGNUMERIC
NUMERIC holds 38 digits with 9 after the point, BIGNUMERIC 76 digits with 38 after the point. A
NUMERIC(P, S) or BIGNUMERIC(P, S) column narrows them further.
Stringis the default and loses nothing. On read it is the canonical text with trailing fractional zeros removed, so0.00comes back as"0". On write it is a plain decimal with an optional sign and no exponent.- Integers are exact both ways. A value with a fractional part read into an integer fails that
row with
OutOfRange. f64works both ways and is lossy, rounded to 9 fractional digits (38 for BIGNUMERIC) on write and parsed into the nearestf64on read.BigQueryDecimal<T>holds any decimal type that round-trips through itsDisplayandFromStr, such asbigdecimal::BigDecimalorrust_decimal::Decimal. The text form is used whatever the type’s own serde does, since some decimal types serialize asf64. For a bare field there areserialize_as_decimalandserialize_as_optional_decimal.
mod bigdecimal { pub type BigDecimal = String; }
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
struct Invoice {
total: BigQueryDecimal<bigdecimal::BigDecimal>,
#[serde(with = "bigquery::serialize_as_optional_decimal")]
discount: Option<bigdecimal::BigDecimal>,
}
On write, a NUMERIC with more than 29 digits before the point or more than 9 after it, a
BIGNUMERIC outside its range or with more than 38 digits after the point, NaN and the infinities
fail that row with OutOfRange. A text that T::from_str
rejects on read, for example rust_decimal::Decimal above its 28 digits, fails with Custom.
Be aware of what BigQuery does on a NUMERIC(P, S) column. The library sends every value at the
full scale of the type, and BigQuery rounds the digits beyond the column’s scale half away from
zero, with no error. So 1.245 and -1.245 written to a NUMERIC(10, 2) column are stored as
1.25 and -1.25. A value over the column’s precision is checked by BigQuery only.
BOOL, STRING and BYTES
BOOL is bool only. Integers and strings are not taken for it.
STRING is String, Box<str>, a char for one character, a unit-variant enum, or any type whose
serde form is a string, such as uuid::Uuid. The maximum length of a STRING(n) column is checked
by BigQuery only.
BYTES is Vec<u8>, serde_bytes::ByteBuf, or [u8; N] when the value has exactly N bytes. A string
is not taken for BYTES on write, and bytes are not taken for STRING: BigQuery stores a string sent
to BYTES as its raw UTF-8 bytes and does not decode base64, so a base64 text would be stored as
text by accident. A serde_json::Value row cannot hold BYTES; select TO_BASE64(b) instead, or use
a typed field.
GEOGRAPHY
GEOGRAPHY is a String. BigQuery returns WKT and takes both WKT and GeoJSON on write, storing
either as a geography. Invalid text is rejected by BigQuery.
Dates and times
The four temporal types map to jiff by default, with no attribute:
| BigQuery type | jiff type | Range | String form | Integer form |
|---|---|---|---|---|
| DATE | jiff::civil::Date | 0001-01-01 to 9999-12-31 | YYYY-MM-DD | i32 days since 1970-01-01 |
| TIME | jiff::civil::Time | 00:00:00 to 23:59:59.999999 | HH:MM:SS[.ffffff] | i64 microseconds of the day |
| DATETIME | jiff::civil::DateTime | 0001-01-01T00:00:00 to 9999-12-31T23:59:59.999999 | YYYY-MM-DDTHH:MM:SS[.ffffff] | i64 civil microseconds since 1970-01-01T00:00:00 |
| TIMESTAMP | jiff::Timestamp | 0001-01-01 to 9999-12-31T23:59:59.999999 UTC | RFC 3339 | i64 microseconds since the epoch |
Some details for each:
- BigQuery keeps microseconds, so sub-microsecond digits are dropped on write, floored for TIMESTAMP;
- a DATE or DATETIME before year 1 on write is
OutOfRange, since jiff allows years down to -9999; - a DATETIME has no offset, so
jiff::Timestampis not taken for it, and a TIMESTAMP is not taken byjiff::civil::DateTimeorjiff::Zoned; - a DATETIME
StringusesTon read and takesTor a space on write; - a TIMESTAMP
Stringis always...Zon read, and on write needsZor±HH:MM, withTor a space between date and time. Other shapes, such as+00orUTC, fail withInvalidText.
Be aware of the TIMESTAMP range. jiff::Timestamp ends at 9999-12-30T22:00:00.999999999Z, a day
before BigQuery’s maximum. A TIMESTAMP above it read into a jiff type, the wrapper included, fails
that row with OutOfRange. String and i64 hold the full range, so read into one of those if
your data can have such values, for example 9999-12-31 used as “never”. The same holds for the
TIMESTAMP ends of a RANGE<TIMESTAMP>.
TIMESTAMP columns with picosecond precision (timestamp_precision = 12) are not supported yet, and
fail with UnsupportedType.
Temporal wrappers
jiff’s serde speaks text only, so a plain jiff field is printed as text by the codec and parsed by jiff, on both read and write. The wrappers keep the jiff types and skip the text:
BigQueryTimestamp(pub jiff::Timestamp);BigQueryDate(pub jiff::civil::Date);BigQueryTime(pub jiff::civil::Time);BigQueryDateTime(pub jiff::civil::DateTime).
With the library’s codecs they read and write BigQuery’s integers directly. In any other serde
format they are the same as the plain jiff type, so a row of wrappers serializes to the same JSON as
a row of plain jiff fields. For a field you do not want to change the type of, the with modules do
the same:
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
struct Event {
happened_at: BigQueryTimestamp,
day: Option<BigQueryDate>,
#[serde(with = "bigquery::serialize_as_timestamp")]
received_at: jiff::Timestamp,
#[serde(with = "bigquery::serialize_as_optional_datetime")]
local_time: Option<jiff::civil::DateTime>,
}
The modules are serialize_as_{timestamp,date,time,datetime} and their
serialize_as_optional_.. forms. A wrapper is checked against its column, so a BigQueryDate on a
TIMESTAMP column fails with TypeMismatch, since its integer form would read microseconds as days.
Use them for wide tables or hot paths. benches/read_codec.rs and benches/write_codec.rs compare
plain jiff fields, the wrappers and integers.
JSON
A JSON column maps by the Rust type, with no wrapper:
- a
Stringfield is the JSON text as it is, in both directions; - any other serde shape, such as
serde_json::Value, a typed struct, a map or a sequence, is parsed on read and printed as JSON on write.
The same holds for a JSON field inside a STRUCT and for the elements of an ARRAY<JSON>:
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
struct Payload {
kind: String,
amount: i64,
}
#[derive(Debug, Serialize, Deserialize)]
struct Message {
// JSON column, as text
raw: String,
// JSON column, as any JSON value
document: serde_json::Value,
// JSON column, parsed into a struct
payload: Option<Payload>,
// ARRAY<JSON> column
events: Vec<serde_json::Value>,
}
JSON has two kinds of null, SQL NULL and the JSON value null, and the rules for them are:
Option<serde_json::Value>keeps them apart: SQL NULL isNone, JSONnullisSome(Value::Null);Option<String>gets SQL NULL asNoneand JSONnullas the text"null";- any other
Option<T>, such asOption<Payload>orOption<i64>, reads both asNone; - a bare typed field, such as
Payload, fails a JSONnullwithCustom.
Text that does not parse into the Rust type on read fails that row with Custom. On write, invalid
JSON text from a String field is checked by BigQuery only: the append fails with RowErrors, with
a FIELDS_ERROR for that row, and no row of that batch is stored.
BigQueryJson<T> is for the two cases the rule above leaves out:
BigQueryJson<String>is a JSON string value parsed into aString, where a plainStringfield would be the raw text with its quotes;- a top-level
Value::Stringis written the way aStringis, as the JSON text itself. A document that can be a bare JSON string needsBigQueryJson<serde_json::Value>.
In other serde formats BigQueryJson<T> is the JSON text as a string. For a bare field there are
serialize_as_json and serialize_as_optional_json.
INTERVAL
An INTERVAL has three parts, each with its own sign, so -1 month +3 days is a valid value. That is
what BigQueryInterval holds:
use bigquery::*;
let interval = BigQueryInterval {
months: -1,
days: 3,
nanos: 4 * 3_600 * 1_000_000_000,
};
// jiff::Span has one sign for all its units, so the conversion can fail
let span: Result<jiff::Span, _> = interval.try_into();
assert!(span.is_err());
let interval = BigQueryInterval::try_from(jiff::Span::new().days(3).hours(4))?;
assert_eq!(interval.days, 3);
Ok::<(), bigquery::errors::BigQueryError>(())
A String works too, in BigQuery’s canonical form [-]Y-M [-]D [-]H:M:S[.ffffff], for example
1-2 -3 4:5:6.000789.
BigQuery keeps microseconds, so nanos that are not whole microseconds fail the row with
OutOfRange on write. Be aware of one BigQuery limit: Storage Read fails the whole stream for an
INTERVAL whose time part is beyond about 2,562,047 hours, and no Arrow client can work around it. So
the library refuses such a value on write, where it could not be read back. A value like that
written another way can still be read with CAST(iv AS STRING) in a query.
RANGE
BigQueryRange<T> holds a RANGE<DATE>, RANGE<DATETIME> or RANGE<TIMESTAMP>, with T any
Rust form of the element type, the temporal wrappers included:
use bigquery::*;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
struct Promotion {
valid: Option<BigQueryRange<jiff::civil::Date>>,
window: BigQueryRange<BigQueryTimestamp>,
}
start is inclusive and end exclusive, and None is an unbounded end. A NULL range into a bare
BigQueryRange<T> fails with NullForNonOption. Whether start comes before end is checked by
BigQuery.
Be aware of a REQUIRED RANGE column: BigQuery reads an unbounded end of it as 1970-01-01 (or the
epoch for DATETIME and TIMESTAMP), and the library cannot tell it from a real epoch bound. NULLABLE
and REPEATED RANGE columns keep unbounded ends as None.
STRUCT
A STRUCT column is a struct with derived serde, its fields matched by name, or a map with string
keys, or a serde_json::Value object. Nested structs and ARRAY<STRUCT> work the same way.
- a NULL STRUCT into a bare struct fails with
NullForNonOption, also when every field of the struct is anOption, so useOption<Address>for a NULLABLE STRUCT; - on write a Rust field with no column fails with
UnknownFieldnaming the path, for exampleshipping.inner.zip, and a REQUIRED subfield the row never wrote withMissingRequiredField; - a tuple works as a whole row only, a STRUCT column into a tuple is
TypeMismatch.
Full example available here.
Field names and serde attributes
Fields are matched to columns by name, so the order of the fields does not matter. The usual serde attributes work in both directions:
renameandrename_all, for example for camelCase columns;alias, for a column that has another name in some tables;skip,default,flattenanddeny_unknown_fields;- a missing
Optioncolumn reads asNone, and a missing column withdefaulttakes its default.
When an aliased field finds both its own column and the alias column in one table, the row fails
with serde’s duplicate field, the same as serde_json gives for that input. A hand-written
Deserialize that answers struct fields by its own integer numbering is not supported; one that
takes field names works.
path! builds column names from Rust fields. For a field renamed to camelCase use
path_camel_case!, for any other #[serde(rename)] the column name as a string.
Dynamic rows
A row can be a serde_json::Value, a map with string keys, or an #[serde(untagged)] enum on read,
when the shape is not known up front. For whole Arrow batches with no serde at all, the read and
query builders have record_batches().
Query parameters
Query parameters take the same Rust forms. .param(name, value) infers the type from the value’s
serde form:
- integers are INT64, floats FLOAT64, strings STRING,
boolBOOL,serde_bytesvalues BYTES; - sequences are ARRAY, structs and string-keyed maps STRUCT;
- the wrappers are their own types:
BigQueryTimestampis TIMESTAMP,BigQueryDateDATE, etc.,BigQueryJson<T>JSON andBigQueryIntervalINTERVAL; BigQueryDecimal<T>is NUMERIC, or BIGNUMERIC for a value NUMERIC cannot hold;BigQueryRange<T>is RANGE when its bounds are temporal wrappers.
A Vec<u8> serializes as a sequence of integers, so it is an ARRAY<INT64>; use
serde_bytes::ByteBuf or param_as for BYTES.
A plain jiff value serializes as text, so it is inferred as a STRING. None, an empty sequence and
a sequence of mixed types cannot be inferred either. For those, .param_as(name, type, value) sets
the type and takes the value in any form the write path takes for it, and None there is a NULL of
that type:
use bigquery::*;
async fn example(db: BigQueryDb, since: jiff::Timestamp) -> BigQueryResult<()> {
let outcome = db
.fluent()
.query(
"SELECT COUNT(*) AS n FROM shop.orders \
WHERE placed_at >= @since AND (@customer IS NULL OR customer = @customer)",
)
.param_as("since", BigQueryFieldType::Timestamp, since)
.param_as(
"customer",
BigQueryFieldType::String { max_length: None },
None::<String>,
)
.execute()
.await?;
let _ = outcome;
Ok(())
}
BigQueryParamType::array_of(..) declares an ARRAY parameter.
Schema types
Table schemas use one vocabulary, whatever API they came from: BigQueryTableSchema with its
BigQueryFieldSchema columns, each with a name, a BigQueryFieldType, a BigQueryFieldMode, a
description and a default value expression.
The v2 API names types with the legacy names, the Storage API with an enum, and Arrow with physical
types only, so the library normalises all of them. INTEGER is Int64, FLOAT is Float64,
BOOLEAN is Bool, RECORD is Struct, DECIMAL is Numeric, BIGDECIMAL is BigNumeric, and
an empty mode, as DDL creates it, is Nullable. Type parameters are part of the type:
String { max_length }, Numeric(Some(BigQueryDecimalParams { precision, scale })), etc. So two
schemas compare with ==.
Display prints GoogleSQL type syntax, such as STRING(10), NUMERIC(10, 2), RANGE<DATE> and
STRUCT<a INT64, b ARRAY<STRING>>.
Codec errors
A value that does not fit fails as BigQueryError::SerializeError on write and
BigQueryError::DeserializeError on read. Both carry a BigQuerySerializationError with a
BigQueryCodecErrorKind, the field path, such as recs[1].v, and on read the row index. A read
error fails only its own row: the other rows of the batch still decode.
| Kind | When |
|---|---|
TypeMismatch | the Rust form is not one the column’s type takes, or a temporal wrapper is on a column of another type |
NullForNonOption | read: NULL into a target that is not an Option, a NULL STRUCT or RANGE included |
NullForRequired | write: None for a REQUIRED field |
NullArrayElement | write: None inside a REPEATED field |
OutOfRange | the value is outside the BigQuery type (a year before 1, NUMERIC digits, the INTERVAL time part) or outside the Rust target (a narrower integer, jiff’s TIMESTAMP maximum) |
InvalidText | a text form that does not parse, such as a malformed DATE, TIMESTAMP, NUMERIC or INTERVAL |
UnknownField | write: a field or map key with no column |
MissingRequiredField | write: a REQUIRED field the row never wrote |
UnsupportedType | a column type the library does not handle, such as TIMESTAMP with picosecond precision |
RowTooLarge | write: one encoded row is larger than the request budget |
Custom | anything the target type’s own serde impl raises, such as missing field or a FromStr failure |
Observability
The library uses tracing. Every call opens one span at the
DEBUG level, and the spans carry what a call used besides how long it took: bytes processed
and billed, slot milliseconds, rows and bytes read or sent, etc.
The spans are:
BigQuery Query: one query terminal call,query(),execute(),dry_run()etc.;BigQuery Read: one table read through the Storage Read API, including the read of a large query result;BigQuery streaming write: one writer, frominsert(),create_streaming_writerorcreate_cdc_writer;BigQuery commit write streams: onecommit_write_streamscall;BigQuery Cancel Job, and the admin spansBigQuery dataset,BigQuery datasets,BigQuery table,BigQuery tables,BigQuery job,BigQuery jobsandBigQuery schema, which carry only the resource they work on.
Retries are logged at WARN inside the span of the call that retried.
Showing the spans
Any tracing subscriber works. With tracing-subscriber, an EnvFilter that enables DEBUG for
the library and span close events, you see every span with its fields when it ends:
use bigquery::*;
use tracing_subscriber::fmt::format::FmtSpan;
async fn example() -> BigQueryResult<()> {
tracing_subscriber::fmt()
.with_env_filter("info,bigquery=debug")
.with_span_events(FmtSpan::CLOSE)
.init();
let db = BigQueryDb::new("my-gcp-project-id").await?;
let outcome = db
.fluent()
.query("SELECT word FROM `bigquery-public-data.samples.shakespeare` LIMIT 10")
.execute()
.await?;
let _ = outcome;
Ok(())
}
Full example available here.
OpenTelemetry
With tracing-opentelemetry every span field
becomes a span attribute with the same name, such as /bigquery/bytes_billed, so you can see the
cost of a request next to its latency in Cloud Trace, Jaeger, etc.
The fields are declared empty when the span opens and recorded once their value is known, which
tracing-opentelemetry exports the same way as fields given at the start. Add its layer to your
subscriber as usual:
use opentelemetry::trace::TracerProvider;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
let provider = opentelemetry_sdk::trace::SdkTracerProvider::builder()
.with_batch_exporter(exporter) // any exporter, e.g. opentelemetry-otlp
.build();
tracing_subscriber::registry()
.with(tracing_subscriber::EnvFilter::new("info,bigquery=debug"))
.with(tracing_opentelemetry::layer().with_tracer(provider.tracer("my-service")))
.init();
Be aware the library spans are DEBUG, so your filter has to enable that level for bigquery.
Common rules for the fields
- A figure BigQuery did not report is left unrecorded, never recorded as
0. The two estimates of a read session are the exception: BigQuery sends them as plain numbers, with no way to tell0from unset. - The library makes no call only for the stats. Every figure comes from a response the call receives anyway, or is counted locally.
- No span field carries the SQL text, a parameter value or a row filter. The query span has the length of the SQL only.
- A row that
stream_query()or a table read skips is logged with the kind of error, the row and the field path, without the message, which can contain the cell’s text. A retry is logged with BigQuery’s error message, as BigQuery wrote it.
Query span
BigQuery Query covers a query from the Query call until the library knows where the rows are.
It does not cover streaming the rows: a result read through Storage Read has its own
BigQuery Read span, a sibling of the query span under your current span.
| Field | What it means | Where it comes from |
|---|---|---|
/bigquery/sql_len | the length of the SQL text in bytes | the SQL, when the span opens |
/bigquery/job_id | the job that ran the query; unset for a query BigQuery ran without a job | the job reference of the Query response |
/bigquery/query_id | the ID BigQuery gave the query, with or without a job | the Query response |
/bigquery/location | where the query ran | the job reference, otherwise the location of the Query response |
/bigquery/statement_type | SELECT, INSERT, UPDATE, CREATE_TABLE, etc. | the Query response, otherwise the job’s query statistics |
/bigquery/bytes_processed | bytes the query processed | the Query response, GetQueryResults, then the job’s statistics |
/bigquery/bytes_billed | bytes billed, after BigQuery’s rounding and minimums | the Query response, then the job’s query statistics |
/bigquery/slot_ms | slot milliseconds the query used | the Query response, then the job’s statistics |
/bigquery/cache_hit | whether the query cache answered | the Query response, GetQueryResults, then the job’s query statistics |
/bigquery/dml_rows | rows a DML statement changed | the same three as cache_hit |
/bigquery/total_rows | rows in the result | the Query response, then GetQueryResults |
/bigquery/route | where the rows came from: inline, storage_read or none | the library, for the terminals that read rows |
Each figure is taken from the first response that reports it, in the order of the table. Which of them a query can have depends on its route:
- Inline, a result complete in the first response: only the
Queryresponse’s own figures, since there is no further call on this route. In the live tests BigQuery reported bytes processed, bytes billed, slot milliseconds and cache hit there; - Storage Read, a larger result: the
Queryresponse, and then the statistics of the job that the library reads anyway to find the destination table; - Polled, a job that was not complete in the first response:
GetQueryResultsand the job’s statistics.GetQueryResultshas no bytes billed and no slot milliseconds, so those come from the job only; execute(): the same as inline when the job completes in the first response.routeis not recorded, sinceexecute()reads no rows;dry_run(): none of the figures. Its estimate is inBigQueryDryRunResultinstead.
The same figures are in BigQueryJobStats from query_with_stats(), stream_query_with_stats()
and execute(), see Job stats.
slot_ms is not 0 even for a query that reads no table. On 1,000 generated rows BigQuery
reported 0 bytes processed and billed, and 25 slot ms inline, 152 slot ms through Storage Read.
Read span
BigQuery Read covers one read session, from opening it until the stream of rows or batches is
dropped.
| Field | What it means | Where it comes from |
|---|---|---|
/bigquery/table | the table being read, project.dataset.table or dataset.table | the read, when the span opens |
/bigquery/streams | the read streams BigQuery gave the session | the read session |
/bigquery/estimated_bytes_scanned | BigQuery’s estimate of the bytes the session scans | the read session, as BigQuery sends it |
/bigquery/estimated_rows | BigQuery’s estimate of the rows the session returns | the read session, as BigQuery sends it |
/bigquery/rows_read | rows received over every stream | the sum of the row counts of every ReadRows response |
/bigquery/bytes_read | uncompressed bytes received | the sum of the uncompressed sizes the ReadRows responses report, when any of them does |
/bigquery/throttle_percent | how much BigQuery throttled the read, 0 to 100 | the highest throttle state any ReadRows response reported, when any of them does |
The three running totals are recorded when the stream is dropped, which also happens when it ends or fails. A read you drop early records what it got so far.
The number of streams is what BigQuery gave, which can be fewer than the library asked for. In the benchmarks the library asked for 16 and got 4 for a 1M-row table.
Write span
BigQuery streaming write covers one writer, from opening its write stream until its background
task ends. insert() and the CDC writer go through the same writer, so they have the same span.
| Field | What it means | Where it comes from |
|---|---|---|
/bigquery/table | the table being written | the writer, when the span opens |
/bigquery/write_mode | Default, Committed or Pending | the writer’s options, when the span opens |
/bigquery/rows_appended | rows in the batches BigQuery acknowledged | counted by the writer, the same as rows_written in BigQueryWriteSummary |
/bigquery/bytes_sent | the encoded size of every AppendRows request sent, resends included, before gRPC framing | counted by the writer, the same as bytes_sent in BigQueryWriteSummary |
/bigquery/appends | AppendRows requests sent, resends included | counted by the writer |
/bigquery/retries | requests that sent a batch again | counted by the writer |
The four counts are recorded when the writer’s background task ends: after finish() or
finalize(), when the writer is dropped, and when it fails for good.
BigQuery commit write streams has only /bigquery/table; the commit time is in the result of the
commit.
Benchmarks
The library was compared with Google’s own clients on the same machine, region and data:
- the official Rust crate google-cloud-bigquery;
- Python
google-cloud-bigquerywithgoogle-cloud-bigquery-storageandpyarrow; - the
bqCLI, for query latency only.
The numbers are from 2026-10-04, both on 0.1.0: the full run of every scenario, and a second run of the query scenarios, the query run. They depend a lot on the network, so treat them as a comparison between the clients on one connection, not as absolute figures.
Method
- Machine: Intel Core i7-10700K, 16 threads, 64 GB RAM, Linux 7.2, CPU governor
powersave. A home connection in Sweden. - Region: every table in
europe-north2(Stockholm), and every query sent with that location, so generated rows are computed there too. - Endpoints: the global defaults for every client,
bigquery.googleapis.comandbigquerystorage.googleapis.com. The TCP connect to both took about 10 ms and the TLS handshake about 29 ms. A small authenticated call,datasets.getover the library’s warm gRPC channel, took about 0.1 s. - Data: generated once with
CREATE TABLE ... AS SELECToverGENERATE_ARRAY, in a scratch dataset deleted at the end. The scan table has 1M rows and 20 columns: INT64, FLOAT64, STRING, BOOL, NUMERIC, DATE and TIMESTAMP, a STRUCT and an ARRAY, some of them nullable. It is about 215 MB as BigQuery counts it. - Runs: 1 warm-up and 5 measured runs for every client and scenario, timed inside the
client process (outside of it for
bq). The tables show the median and the range. - Order: one client at a time. The clients run round by round, run 1 of every client, then run 2, etc., with the order rotated each round, so drifts in the machine or the network spread over all of them.
- Quiet machine: before each run the harness waited for a 1-minute load under 1.6 (10% of the cores), no other process above 20% of a core, no compiler or linker running and under 2 MB/s of idle network traffic. During the run it sampled the machine every 2 s and repeated a run with a bad sample. In the full run 4 runs were repeated, all because of a busy editor on the same desktop; in the query run 4 more, 3 of them because of a Rust build on the same machine. None were left contended. Every reported run had other processes using under 0.6 cores, and the network traffic during the runs matched what the clients themselves moved.
- Settings: the defaults of each client, query cache off everywhere. Where a client leaves a
setting to the caller, the harness set it:
- Storage Read on the official crate: 16 streams asked for (the library asks for the machine’s parallelism), LZ4 buffers, one task per stream;
- Storage Write on the official crate: 25,000 rows per Arrow request, about 5.5 MB, so a request stays under the 10 MB limit, and 8 requests in flight, which is the library’s default window.
Python asks Storage Read for max_stream_count = 0 by default, which lets BigQuery decide. It got
1 stream every time, while on the 1M-row table scan the library and the official crate got 4 of the
16 they asked for.
Versions
| Client | Versions |
|---|---|
| bigquery (this library) | 0.1.0; gcloud-sdk 0.32.4 (0.32.3 for the full run), tonic 0.14.6, arrow 60.0.0 |
| google-cloud-bigquery | 0.18.0, google-cloud-bigquery-v2 1.0.0, google-cloud-gax 1.15.0, google-cloud-auth 1.17.0 |
| Python | Python 3.14.7, google-cloud-bigquery 3.46.1, google-cloud-bigquery-storage 2.42.0, pyarrow 25.0.1, pandas 3.0.6, grpcio 1.84.0 |
| bq | BigQuery CLI 2.1.39 (Google Cloud SDK 587.0.0) |
Rust code was built in release mode with rustc 1.98.1.
What each client does
The paths below were checked in each client’s source for the versions above, and the harness records which one each run actually took:
- This library: everything over gRPC. A query goes through the v2
Querycall with an Arrow result andJOB_CREATION_OPTIONALby default (.job_creation_required()turns it off); a result that does not fit in the first response is read through Storage Read. Typed rows are decoded from Arrow with serde. - google-cloud-bigquery: queries over REST (
google-cloud-bigquery-v2), rows as JSON pages converted with#[derive(FromRow)]. It sends queries withJOB_CREATION_OPTIONALby default, so a short query may run without a job. Storage Read is a raw generated client that returns Arrow IPC bytes and leaves the decoding to you; in 0.18 it is not behind anycfgflag. Storage Write takes Arrow you serialize yourself; its protobuf writer is not public. - Python:
query_and_wait()over REST, andto_arrow()orto_dataframe()through Storage Read for large results and table scans. - bq: one CLI process per query.
Small query latency
The library is measured twice, in its default short query mode and with
.job_creation_required(), so you can see what a job costs. Both run in the same harness run
as the other clients.
SELECT 1 AS x:
| Client | Median | Range |
|---|---|---|
| bigquery | 0.096 s | 0.089-0.115 s |
bigquery, .job_creation_required() | 0.167 s | 0.160-0.256 s |
| google-cloud-bigquery | 0.096 s | 0.086-0.110 s |
| Python | 0.173 s | 0.144-0.272 s |
| bq | 2.354 s | 2.325-2.476 s |
1,000 generated rows:
| Client | Median | Range | Rows/s |
|---|---|---|---|
| bigquery | 0.130 s | 0.110-0.145 s | 7,676 |
bigquery, .job_creation_required() | 0.203 s | 0.177-0.207 s | 4,930 |
| google-cloud-bigquery | 0.134 s | 0.117-0.137 s | 7,480 |
| Python | 0.213 s | 0.195-0.284 s | 4,690 |
| bq | 2.392 s | 2.338-2.402 s | 418 |
The library and the official crate are the same speed here, and both are faster than Python, because of the job:
- The official crate sets
JobCreationMode::JobCreationOptionalon every query (google-cloud-bigquery0.18.0,src/query/builder.rs, line 73), BigQuery’s short query mode, where BigQuery may answer without creating a job. The library does the same by default. In these runs neither got a job back, while Python and the library’s required mode got one every time. - Measured in one process on the library’s own channel, the library adds nothing over a raw
Querycall: 0.086 s against 0.085-0.092 s for the constant query in short query mode, and 0.151 s against 0.149 s with a job. Creating the job is the whole difference, about 65 ms.
About 0.6 s of every bq run is the CLI starting up (bq version timed in each run); the rest
probably is its job insert and polling, I didn’t check.
Large query result
200,000 generated rows of 4 columns, read as typed rows, from the query run:
| Client | Path | Median | Range | Rows/s |
|---|---|---|---|---|
| bigquery | Storage Read, 1 stream | 1.293 s | 1.246-1.484 s | 154,674 |
bigquery, .job_creation_required() | Storage Read, 1 stream | 1.279 s | 1.100-1.564 s | 156,343 |
| google-cloud-bigquery | REST JSON pages | 3.862 s | 3.617-4.088 s | 51,787 |
Python (list(rows)) | REST JSON pages | 4.704 s | 4.567-5.004 s | 42,515 |
A result this large always gets a job, so the two modes are the same here. The full run gave 1.260 s, 3.690 s and 4.935 s.
The same result as Arrow, from the full run:
| Client | Path | Median | Range | Rows/s |
|---|---|---|---|---|
bigquery (record_batches()) | Storage Read, 1 stream | 1.124 s | 1.107-1.202 s | 177,868 |
Python (to_arrow()) | Storage Read, 1 stream | 2.373 s | 2.351-2.608 s | 84,289 |
| google-cloud-bigquery | n/a: its query client returns JSON rows only |
The library reads the destination table through Storage Read once the result does not fit in
the first response, while the official crate and Python iterate REST pages. The official crate
spent 0.019 s of its 3.7-3.9 s in FromRow, so nearly all of its time is the paging itself.
Python takes the same Storage Read path for to_arrow() and is still twice as slow. I guess it
is the extra metadata calls before the read session, but I didn’t measure that.
Table scan
The whole 1M-row table, every column. MB/s is the table’s 215 MB divided by the median.
Typed rows:
| Client | Path | Median | Range | Rows/s | MB/s |
|---|---|---|---|---|---|
bigquery (obj::<T>()) | Storage Read, 4 streams | 9.764 s | 7.708-10.307 s | 102,414 | 22.0 |
Python (to_dataframe()) | Storage Read, 1 stream | 10.746 s | 10.494-11.396 s | 93,061 | 20.0 |
google-cloud-bigquery (FromRow) | SELECT * query, REST JSON pages | 94.723 s | 92.798-94.884 s | 10,557 | 2.3 |
Arrow:
| Client | Path | Median | Range | Rows/s | MB/s |
|---|---|---|---|---|---|
bigquery (record_batches()) | Storage Read, 4 streams | 9.344 s | 7.516-9.722 s | 107,018 | 23.0 |
google-cloud-bigquery (raw client + arrow-ipc) | Storage Read, 4 streams | 9.276 s | 7.336-9.713 s | 107,801 | 23.1 |
Python (to_arrow()) | Storage Read, 1 stream | 9.918 s | 9.767-10.579 s | 100,830 | 21.6 |
The official crate has no typed Storage Read, so its typed path is a query, and a SELECT *
query bills the table: 215 MB per run. Its FromRow took about 1.05 s of the 95 s.
Raw Arrow scans are the same speed in all three clients, and the library has no advantage there. Every client moved about 150 MB per scan on the network interface at about 16 MB/s, so I think this connection is the limit, not the clients.
Python’s to_dataframe() builds a pandas frame with typed columns, not row objects, so it is the
closest Python has to typed rows rather than the same thing.
Storage Write
1M rows of the scan table’s shape into the default stream, per run:
| Client | Format | Median | Range | Rows/s | MB sent |
|---|---|---|---|---|---|
bigquery (insert().objects()) | protobuf from serde | 29.957 s | 25.186-30.267 s | 33,381 | 180.5 |
| google-cloud-bigquery | Arrow built from the same rows | 32.715 s | 30.786-33.363 s | 30,567 | 221.4 |
| Python | n/a: its writer takes requests you build yourself, protobuf rows with a hand-made descriptor |
Both times include turning the Rust rows into the wire format. The library sent 18% fewer bytes. On the network interface the library ran at about 6.5 MB/s and the official crate at about 7.3 MB/s, so the upload of this connection is probably not the whole limit. I think the smaller requests are why the library is a bit faster, but that is not measured.
Decode cost
The scan table’s 1M rows already in memory as Arrow batches, decoded into structs on one thread:
| Client | Median | Range | Rows/s |
|---|---|---|---|
| bigquery | 0.649 s | 0.642-0.651 s | 1,541,720 |
| google-cloud-bigquery | n/a |
The official crate has no Arrow to struct decoder. Its FromRow converts the JSON rows of a live
query, and Row has no public constructor, so it cannot be fed the same batches. Inside the REST
scan above its conversion took about 1.05 s per 1M rows, but on already parsed JSON values, so the
two numbers are not comparable.
In the library the decode runs on each read stream’s task, next to the network reads, which is probably why typed and Arrow scans take almost the same time.
Cost
The full run billed 1.72 GB of queries, all of them the official crate’s SELECT * scans, and the
query run billed 0 bytes. Every other query read generated rows and billed 0 bytes. The scans
went through the Storage Read free tier, and the writes ingested about 215 MB per 1M rows under the Storage Write free tier.
Reproducing
The harness is in bench-compare,
an unpublished crate with the Rust contenders and a uv project for
Python. You need application default credentials and optionally the bq CLI:
bench-compare/run.sh --project my-project
It creates the scratch dataset in europe-north2 (--location changes it), runs everything and
deletes the dataset, also on failure. The raw results land in
bench-compare/results/<run label>/results.json, with the machine state of every run. A full run
takes about 35 minutes and costs a few cents. --only query_const,query_1k,query_200k_rows runs
just the query scenarios. The summary of both runs above is in
bench-compare/results-2026-10-04-europe-north2.json.