google_cloud_bigquery/write/
client.rs1use super::client_builder::ClientBuilder;
16use super::pool::{StreamPool, StreamPoolOptions};
17use super::retry_policy::RetryOptions;
18use super::stream_type::{ApplicationCreatedStream, DefaultStream};
19use super::transport::Transport;
20use super::writer_builder::WriterBuilder;
21use crate::ClientBuilderResult as BuilderResult;
22use std::collections::HashMap;
23use std::sync::{Arc, Mutex};
24
25#[derive(Debug)]
27pub struct Write {
28 inner: Arc<Transport>,
29 pools: Arc<Mutex<HashMap<String, Arc<StreamPool>>>>,
30 pool_options: StreamPoolOptions,
31 retry_options: RetryOptions,
32}
33
34impl Write {
35 pub fn builder() -> ClientBuilder {
37 ClientBuilder::new()
38 }
39
40 pub(crate) async fn new(builder: ClientBuilder) -> BuilderResult<Self> {
41 let inner = Arc::new(Transport::new(builder.config).await?);
42 let pools = Arc::new(Mutex::new(HashMap::new()));
43 Ok(Self {
44 inner,
45 pools,
46 pool_options: builder.pool_options,
47 retry_options: builder.retry_options,
48 })
49 }
50
51 pub fn open_default_stream<T: Into<String>>(&self, table: T) -> WriterBuilder<DefaultStream> {
74 WriterBuilder::new_open_default(
75 self.inner.clone(),
76 self.pools.clone(),
77 self.pool_options.clone(),
78 self.retry_options.clone(),
79 table.into(),
80 )
81 }
82
83 pub fn create_stream<S: ApplicationCreatedStream, T: Into<String>>(
119 &self,
120 table: T,
121 ) -> WriterBuilder<S> {
122 WriterBuilder::new_create(self.inner.clone(), self.retry_options.clone(), table.into())
123 }
124
125 pub fn attach_to_stream<S: ApplicationCreatedStream, T: Into<String>>(
158 &self,
159 write_stream: T,
160 ) -> WriterBuilder<S> {
161 WriterBuilder::new_attach(
162 self.inner.clone(),
163 self.retry_options.clone(),
164 write_stream.into(),
165 )
166 }
167}
168
169#[cfg(test)]
170mod tests {
171 use super::super::error::AppendError;
172 use super::*;
173 use crate::model::{ArrowRecordBatch, ArrowSchema, ProtoRows, ProtoSchema};
174 use bigquery_grpc_mock::{MockBigQueryWrite, start};
175 use gaxi::grpc::tonic::Status as TonicStatus;
176 use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
177
178 #[tokio::test]
179 async fn arrow() -> anyhow::Result<()> {
180 let mut mock = MockBigQueryWrite::new();
181 mock.expect_append_rows()
182 .return_once(|_| Err(TonicStatus::failed_precondition("fail")));
183 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
184 let client = Write::builder()
185 .with_endpoint(endpoint)
186 .with_credentials(Anonymous::new().build())
187 .build()
188 .await?;
189 let writer = client
190 .open_default_stream("projects/p/datasets/d/tables/t")
191 .build_arrow(ArrowSchema::new())
192 .await?;
193 let err = writer
194 .append(ArrowRecordBatch::new())
195 .send()
196 .await
197 .expect_err("write should fail");
198 assert!(matches!(err, AppendError::Rpc { source: _ }));
199
200 Ok(())
201 }
202
203 #[tokio::test]
204 async fn proto() -> anyhow::Result<()> {
205 let mut mock = MockBigQueryWrite::new();
206 mock.expect_append_rows()
207 .return_once(|_| Err(TonicStatus::failed_precondition("fail")));
208 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
209 let client = Write::builder()
210 .with_endpoint(endpoint)
211 .with_credentials(Anonymous::new().build())
212 .build()
213 .await?;
214 let writer = client
215 .open_default_stream("projects/p/datasets/d/tables/t")
216 .build_proto(ProtoSchema::new())
217 .await?;
218 let err = writer
219 .append(ProtoRows::new())
220 .send()
221 .await
222 .expect_err("write should fail");
223 assert!(matches!(err, AppendError::Rpc { source: _ }));
224
225 Ok(())
226 }
227
228 #[tokio::test]
229 async fn multiplexing() -> anyhow::Result<()> {
230 let mut mock = MockBigQueryWrite::new();
231 mock.expect_get_write_stream().times(2).returning(|req| {
232 let name = req.into_inner().name;
233 Ok(gaxi::grpc::tonic::Response::new(
234 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::WriteStream {
235 name,
236 location: "us".to_string(),
237 ..Default::default()
238 },
239 ))
240 });
241 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
242 let client = Write::builder()
243 .with_endpoint(endpoint)
244 .with_credentials(Anonymous::new().build())
245 .build()
246 .await?;
247 let multiplexed_writer1 = client
248 .open_default_stream("projects/p/datasets/d/tables/t1")
249 .with_multiplexing(true)
250 .build_arrow(ArrowSchema::new())
251 .await?;
252 let multiplexed_writer2 = client
253 .open_default_stream("projects/p/datasets/d/tables/t2")
254 .with_multiplexing(true)
255 .build_arrow(ArrowSchema::new())
256 .await?;
257 assert!(Arc::ptr_eq(
258 &multiplexed_writer1.inner.pool,
259 &multiplexed_writer2.inner.pool
260 ));
261
262 let standalone_writer = client
263 .open_default_stream("projects/p/datasets/d/tables/t3")
264 .with_multiplexing(false)
265 .build_arrow(ArrowSchema::new())
266 .await?;
267 assert!(!Arc::ptr_eq(
268 &multiplexed_writer1.inner.pool,
269 &standalone_writer.inner.pool
270 ));
271
272 Ok(())
273 }
274
275 #[tokio::test]
276 async fn format_isolation() -> anyhow::Result<()> {
277 let mut mock = MockBigQueryWrite::new();
278 mock.expect_get_write_stream().times(2).returning(|req| {
279 let name = req.into_inner().name;
280 Ok(gaxi::grpc::tonic::Response::new(
281 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::WriteStream {
282 name,
283 location: "us".to_string(),
284 ..Default::default()
285 },
286 ))
287 });
288 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
289 let client = Write::builder()
290 .with_endpoint(endpoint)
291 .with_credentials(Anonymous::new().build())
292 .build()
293 .await?;
294 let arrow_writer = client
295 .open_default_stream("projects/p/datasets/d/tables/t1")
296 .with_multiplexing(true)
297 .build_arrow(ArrowSchema::new())
298 .await?;
299 let proto_writer = client
300 .open_default_stream("projects/p/datasets/d/tables/t2")
301 .with_multiplexing(true)
302 .build_proto(ProtoSchema::new())
303 .await?;
304
305 assert!(!Arc::ptr_eq(
307 &arrow_writer.inner.pool,
308 &proto_writer.inner.pool
309 ));
310
311 Ok(())
312 }
313
314 #[tokio::test]
315 async fn location_isolation() -> anyhow::Result<()> {
316 let mut mock = MockBigQueryWrite::new();
317 mock.expect_get_write_stream().times(3).returning(|req| {
318 let name = req.into_inner().name;
319 let location = if name.contains("t1") {
320 "us".to_string()
321 } else if name.contains("t2") {
322 "eu".to_string()
323 } else {
324 "us".to_string()
325 };
326 Ok(gaxi::grpc::tonic::Response::new(
327 bigquery_grpc_mock::google::cloud::bigquery::storage::v1::WriteStream {
328 name,
329 location,
330 ..Default::default()
331 },
332 ))
333 });
334 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
335 let client = Write::builder()
336 .with_endpoint(endpoint)
337 .with_credentials(Anonymous::new().build())
338 .build()
339 .await?;
340 let us_writer1 = client
341 .open_default_stream("projects/p/datasets/d/tables/t1")
342 .with_multiplexing(true)
343 .build_arrow(ArrowSchema::new())
344 .await?;
345 let eu_writer = client
346 .open_default_stream("projects/p/datasets/d/tables/t2")
347 .with_multiplexing(true)
348 .build_arrow(ArrowSchema::new())
349 .await?;
350
351 assert!(!Arc::ptr_eq(&us_writer1.inner.pool, &eu_writer.inner.pool));
353
354 let us_writer2 = client
355 .open_default_stream("projects/p/datasets/d/tables/t3")
356 .with_multiplexing(true)
357 .build_arrow(ArrowSchema::new())
358 .await?;
359
360 assert!(Arc::ptr_eq(&us_writer1.inner.pool, &us_writer2.inner.pool));
362
363 Ok(())
364 }
365
366 #[tokio::test]
367 async fn retry_options() -> anyhow::Result<()> {
368 let client = Write::builder()
369 .with_credentials(Anonymous::new().build())
370 .build()
371 .await?;
372 let writer = client
373 .open_default_stream("projects/p/datasets/d/tables/t")
374 .build_arrow(ArrowSchema::new())
375 .await?;
376
377 let options = &writer.inner.options;
379 assert!(Arc::ptr_eq(
380 &client.retry_options.retry_policy,
381 &options.retry_policy
382 ));
383 assert!(Arc::ptr_eq(
384 &client.retry_options.backoff_policy,
385 &options.backoff_policy
386 ));
387 assert_eq!(
388 client.retry_options.attempt_timeout,
389 options.attempt_timeout
390 );
391
392 Ok(())
393 }
394}