Skip to main content

faucet_source_bigquery/
storage_read.rs

1//! BigQuery **Storage Read API** (gRPC) Arrow path (#380).
2//!
3//! Reads a table directly as Arrow `RecordBatch`es via the
4//! `google.cloud.bigquery.storage.v1` gRPC service — no `jobs.query`, no
5//! per-row JSON. Used both by the columnar fast path
6//! ([`BigQuerySource::stream_batches`](crate::stream::BigQuerySource)) and, when
7//! the sink is not columnar, by the row path (Arrow → JSON) so a `read_api`
8//! source works with any sink.
9//!
10//! The `gcloud-*` stack here (gRPC over tonic 0.14, auth via `gcloud-auth`) is
11//! a deliberately separate family from the REST `gcp-bigquery-client` used for
12//! the query path; it is only compiled with the `arrow` feature.
13
14use crate::config::BigQuerySourceConfig;
15use crate::stream::BigQuerySource;
16use arrow::array::RecordBatch;
17use faucet_core::columnar::ColumnarPage;
18use faucet_core::{FaucetError, Stream, StreamPage};
19use futures::StreamExt;
20use gcloud_gax::conn::{ConnectionManager, ConnectionOptions, Environment};
21use gcloud_googleapis::cloud::bigquery::storage::v1::big_query_read_client::BigQueryReadClient;
22use gcloud_googleapis::cloud::bigquery::storage::v1::read_rows_response::{Rows, Schema};
23use gcloud_googleapis::cloud::bigquery::storage::v1::read_session::TableReadOptions;
24use gcloud_googleapis::cloud::bigquery::storage::v1::{
25    CreateReadSessionRequest, DataFormat, ReadRowsRequest, ReadSession,
26};
27use serde_json::Value;
28use std::pin::Pin;
29
30use faucet_common_bigquery::BigQueryCredentials;
31
32const STORAGE_DOMAIN: &str = "bigquerystorage.googleapis.com";
33const STORAGE_AUDIENCE: &str = "https://bigquerystorage.googleapis.com/";
34const STORAGE_SCOPES: [&str; 2] = [
35    "https://www.googleapis.com/auth/bigquery.readonly",
36    "https://www.googleapis.com/auth/cloud-platform",
37];
38/// Bump the client's decode cap well above the 4 MiB default — Arrow batches
39/// from the Storage Read API can be large.
40const MAX_DECODE_BYTES: usize = 1 << 30; // 1 GiB
41
42/// Resolve `dataset.table` / `project.dataset.table` into the Storage Read API
43/// resource name `projects/{p}/datasets/{d}/tables/{t}`. Pure.
44pub fn resolve_table(project_id: &str, read_table: Option<&str>) -> Result<String, FaucetError> {
45    let t = read_table.ok_or_else(|| {
46        FaucetError::Config(
47            "BigQuery read_api requires `read_table` (dataset.table or project.dataset.table)"
48                .into(),
49        )
50    })?;
51    match t.split('.').collect::<Vec<_>>().as_slice() {
52        [dataset, table] => Ok(format!(
53            "projects/{project_id}/datasets/{dataset}/tables/{table}"
54        )),
55        [project, dataset, table] => Ok(format!(
56            "projects/{project}/datasets/{dataset}/tables/{table}"
57        )),
58        _ => Err(FaucetError::Config(format!(
59            "BigQuery read_table '{t}' must be 'dataset.table' or 'project.dataset.table'"
60        ))),
61    }
62}
63
64/// Decode one Storage Read API Arrow message. The API sends the IPC schema
65/// once (first response) and each batch as a standalone IPC record-batch
66/// message; concatenating schema + batch bytes yields a decodable IPC stream.
67fn decode_arrow(schema: &[u8], batch: &[u8]) -> Result<Vec<RecordBatch>, FaucetError> {
68    use arrow::ipc::reader::StreamReader;
69    let mut buf = Vec::with_capacity(schema.len() + batch.len());
70    buf.extend_from_slice(schema);
71    buf.extend_from_slice(batch);
72    let reader = StreamReader::try_new(std::io::Cursor::new(buf), None)
73        .map_err(|e| FaucetError::Source(format!("BigQuery Storage Read arrow decode: {e}")))?;
74    let mut out = Vec::new();
75    for b in reader {
76        out.push(
77            b.map_err(|e| FaucetError::Source(format!("BigQuery Storage Read arrow batch: {e}")))?,
78        );
79    }
80    Ok(out)
81}
82
83/// Build the gRPC auth `Environment` from the connector's BigQuery credentials.
84async fn build_environment(auth: &BigQueryCredentials) -> Result<Environment, FaucetError> {
85    use gcloud_auth::credentials::CredentialsFile;
86    use gcloud_auth::token::DefaultTokenSourceProvider;
87
88    let cfg = gcloud_auth::project::Config::default()
89        .with_audience(STORAGE_AUDIENCE)
90        .with_scopes(&STORAGE_SCOPES);
91    let tsp = match auth {
92        BigQueryCredentials::ApplicationDefault => DefaultTokenSourceProvider::new(cfg)
93            .await
94            .map_err(|e| FaucetError::Auth(format!("BigQuery Storage Read ADC auth: {e}")))?,
95        BigQueryCredentials::ServiceAccountKeyPath { path } => {
96            let cf = CredentialsFile::new_from_file(path.clone())
97                .await
98                .map_err(|e| FaucetError::Auth(format!("BigQuery Storage Read key file: {e}")))?;
99            DefaultTokenSourceProvider::new_with_credentials(cfg, Box::new(cf))
100                .await
101                .map_err(|e| FaucetError::Auth(format!("BigQuery Storage Read key file: {e}")))?
102        }
103        BigQueryCredentials::ServiceAccountKey { json } => {
104            let cf = CredentialsFile::new_from_str(json)
105                .await
106                .map_err(|e| FaucetError::Auth(format!("BigQuery Storage Read inline key: {e}")))?;
107            DefaultTokenSourceProvider::new_with_credentials(cfg, Box::new(cf))
108                .await
109                .map_err(|e| FaucetError::Auth(format!("BigQuery Storage Read inline key: {e}")))?
110        }
111    };
112    Ok(Environment::GoogleCloud(Box::new(tsp)))
113}
114
115/// Open a read session for the configured table and stream its Arrow batches.
116fn read_batches(
117    cfg: &BigQuerySourceConfig,
118) -> impl Stream<Item = Result<RecordBatch, FaucetError>> + Send + '_ {
119    async_stream::try_stream! {
120        let env = build_environment(&cfg.auth).await?;
121        let cm = ConnectionManager::new(
122            1,
123            STORAGE_DOMAIN,
124            STORAGE_AUDIENCE,
125            &env,
126            &ConnectionOptions::default(),
127        )
128        .await
129        .map_err(|e| FaucetError::Source(format!("BigQuery Storage Read connect: {e}")))?;
130        let mut client = BigQueryReadClient::new(cm.conn()).max_decoding_message_size(MAX_DECODE_BYTES);
131
132        let table = resolve_table(&cfg.project_id, cfg.read_table.as_deref())?;
133        let read_options = TableReadOptions {
134            selected_fields: cfg.selected_fields.clone(),
135            row_restriction: cfg.row_restriction.clone().unwrap_or_default(),
136            ..Default::default()
137        };
138        let session = ReadSession {
139            data_format: DataFormat::Arrow as i32,
140            table,
141            read_options: Some(read_options),
142            ..Default::default()
143        };
144        let request = CreateReadSessionRequest {
145            parent: format!("projects/{}", cfg.project_id),
146            read_session: Some(session),
147            max_stream_count: cfg.max_streams.max(1),
148            ..Default::default()
149        };
150        let created = client
151            .create_read_session(request)
152            .await
153            .map_err(|e| FaucetError::Source(format!("BigQuery CreateReadSession failed: {e}")))?
154            .into_inner();
155
156        let mut schema_bytes: Option<Vec<u8>> = None;
157        for stream in &created.streams {
158            let rr = ReadRowsRequest { read_stream: stream.name.clone(), offset: 0 };
159            let mut responses = client
160                .read_rows(rr)
161                .await
162                .map_err(|e| FaucetError::Source(format!("BigQuery ReadRows failed: {e}")))?
163                .into_inner();
164            while let Some(msg) = responses
165                .message()
166                .await
167                .map_err(|e| FaucetError::Source(format!("BigQuery ReadRows stream error: {e}")))?
168            {
169                if let Some(Schema::ArrowSchema(s)) = msg.schema {
170                    schema_bytes = Some(s.serialized_schema);
171                }
172                if let Some(Rows::ArrowRecordBatch(rb)) = msg.rows {
173                    let sch = schema_bytes.as_deref().ok_or_else(|| {
174                        FaucetError::Source(
175                            "BigQuery Storage Read: record batch arrived before the Arrow schema"
176                                .into(),
177                        )
178                    })?;
179                    for batch in decode_arrow(sch, &rb.serialized_record_batch)? {
180                        if batch.num_rows() > 0 {
181                            yield batch;
182                        }
183                    }
184                }
185            }
186        }
187    }
188}
189
190/// Columnar fast path: yield one [`ColumnarPage`] per Arrow batch.
191pub fn stream_batches_arrow(
192    src: &BigQuerySource,
193) -> Pin<Box<dyn Stream<Item = Result<ColumnarPage, FaucetError>> + Send + '_>> {
194    let cfg = src.config();
195    Box::pin(async_stream::try_stream! {
196        let inner = read_batches(cfg);
197        futures::pin_mut!(inner);
198        while let Some(batch) = inner.next().await {
199            yield ColumnarPage::new(batch?, None);
200        }
201        tracing::info!(table = ?cfg.read_table, "BigQuery Storage Read columnar stream complete");
202    })
203}
204
205/// Row path for a non-columnar sink: decode Arrow batches to JSON and re-frame
206/// into [`StreamPage`]s of `batch_size` records.
207pub fn stream_pages_arrow(
208    src: &BigQuerySource,
209) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + '_>> {
210    let cfg = src.config();
211    let batch_size = cfg.batch_size;
212    Box::pin(async_stream::try_stream! {
213        let chunk = if batch_size == 0 { usize::MAX } else { batch_size };
214        let mut buffer: Vec<Value> = Vec::new();
215        let inner = read_batches(cfg);
216        futures::pin_mut!(inner);
217        while let Some(batch) = inner.next().await {
218            let batch = batch?;
219            for v in faucet_core::columnar::record_batch_to_values(&batch)? {
220                buffer.push(v);
221                if buffer.len() >= chunk {
222                    let page = std::mem::replace(&mut buffer, Vec::with_capacity(chunk));
223                    yield StreamPage { records: page, bookmark: None };
224                }
225            }
226        }
227        if !buffer.is_empty() {
228            yield StreamPage { records: buffer, bookmark: None };
229        }
230        tracing::info!(table = ?cfg.read_table, "BigQuery Storage Read row stream complete");
231    })
232}
233
234#[cfg(test)]
235mod tests {
236    use super::*;
237
238    #[test]
239    fn resolve_table_two_and_three_part() {
240        assert_eq!(
241            resolve_table("billing-proj", Some("ds.events")).unwrap(),
242            "projects/billing-proj/datasets/ds/tables/events"
243        );
244        assert_eq!(
245            resolve_table("billing-proj", Some("other-proj.ds.events")).unwrap(),
246            "projects/other-proj/datasets/ds/tables/events"
247        );
248    }
249
250    #[test]
251    fn resolve_table_requires_table_and_valid_shape() {
252        assert!(resolve_table("p", None).is_err());
253        assert!(resolve_table("p", Some("just_a_name")).is_err());
254        assert!(resolve_table("p", Some("a.b.c.d")).is_err());
255    }
256}