Skip to main content

google_cloud_bigquery/write/
client.rs

1// Copyright 2026 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     https://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// A client for BigQuery Storage Write API.
26#[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    /// Creates a new [ClientBuilder].
36    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    /// Opens the [default stream] for the given table.
52    ///
53    /// `table` must have the format
54    /// `projects/{project}/datasets/{dataset}/tables/{table}`.
55    ///
56    /// # Example
57    /// ```
58    /// # use google_cloud_bigquery::client::Write;
59    /// # async fn sample(client: Write) -> anyhow::Result<()> {
60    /// let writer = client
61    ///     .open_default_stream("projects/my-project/datasets/my-dataset/tables/my-table")
62    ///     .build_arrow(schema())
63    ///     .await?;
64    /// # Ok(()) }
65    ///
66    /// use google_cloud_bigquery::model::ArrowSchema;
67    /// fn schema() -> ArrowSchema {
68    ///   todo!("Define your table's schema...")
69    /// }
70    /// ```
71    ///
72    /// [default stream]: https://docs.cloud.google.com/bigquery/docs/write-api#default_stream
73    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    /// Creates a new [application-created stream] for the given table.
84    ///
85    /// `table` must have the format
86    /// `projects/{project}/datasets/{dataset}/tables/{table}`.
87    ///
88    /// The stream type `S` can be inferred from the variable's writer type
89    /// annotation
90    /// ([`PendingWriter`][crate::write::PendingWriter],
91    /// [`CommittedWriter`][crate::write::CommittedWriter], or
92    /// [`BufferedWriter`][crate::write::BufferedWriter]) or specified explicitly via turbofish
93    /// (`create_stream::<PendingStream, _>(...)`).
94    ///
95    /// See [Selecting a type] for guidance on choosing a stream type for your
96    /// workload.
97    ///
98    /// # Example
99    /// ```
100    /// use google_cloud_bigquery::write::PendingWriter;
101    /// use google_cloud_bigquery::write::format::Arrow;
102    /// # use google_cloud_bigquery::client::Write;
103    /// # async fn sample(client: Write) -> anyhow::Result<()> {
104    /// let writer: PendingWriter<Arrow> = client
105    ///     .create_stream("projects/my-project/datasets/my-dataset/tables/my-table")
106    ///     .build_arrow(schema())
107    ///     .await?;
108    /// # Ok(()) }
109    ///
110    /// use google_cloud_bigquery::model::ArrowSchema;
111    /// fn schema() -> ArrowSchema {
112    ///   todo!("Define your table's schema...")
113    /// }
114    /// ```
115    ///
116    /// [application-created stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#application-created_streams
117    /// [Selecting a type]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#selecting_a_type
118    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    /// Attaches a writer to an existing [application-created stream].
126    ///
127    /// `write_stream` must have the format
128    /// `projects/{project}/datasets/{dataset}/tables/{table}/streams/{stream}`.
129    ///
130    /// The stream type `S` can be inferred from the variable's writer type
131    /// annotation
132    /// ([`PendingWriter`][crate::write::PendingWriter],
133    /// [`CommittedWriter`][crate::write::CommittedWriter], or
134    /// [`BufferedWriter`][crate::write::BufferedWriter]) or specified explicitly via turbofish
135    /// (`attach_to_stream::<PendingStream, _>(...)`).
136    ///
137    /// # Example
138    /// ```
139    /// use google_cloud_bigquery::write::CommittedWriter;
140    /// use google_cloud_bigquery::write::format::Arrow;
141    /// # use google_cloud_bigquery::client::Write;
142    /// # async fn sample(client: Write) -> anyhow::Result<()> {
143    /// let writer: CommittedWriter<Arrow> = client
144    ///     .attach_to_stream("projects/my-project/datasets/my_dataset/tables/my_table/streams/my_stream")
145    ///     .build_arrow(schema())
146    ///     .await?;
147    /// # Ok(())
148    /// # }
149    /// #
150    /// # use google_cloud_bigquery::model::ArrowSchema;
151    /// # fn schema() -> ArrowSchema {
152    /// #   todo!("Define your table's schema...")
153    /// # }
154    /// ```
155    ///
156    /// [application-created stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#application-created_streams
157    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        // Different formats receive distinct connection pools.
306        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        // Different locations receive distinct connection pools.
352        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        // Same location shares the connection pool.
361        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        // The writer uses the client's policies, not a fresh set of defaults.
378        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}