faucet_source_bigquery/
storage_read.rs1use 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];
38const MAX_DECODE_BYTES: usize = 1 << 30; pub 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
64fn 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
83async 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
115fn 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
190pub 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
205pub 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}