Skip to main content

google_cloud_bigquery/write/
client_builder.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::pool::StreamPoolOptions;
16use super::retry_policy::RetryOptions;
17use crate::ClientBuilderResult as BuilderResult;
18use crate::client::Write;
19use gaxi::options::ClientConfig;
20use google_cloud_auth::credentials::Credentials;
21use google_cloud_gax::backoff_policy::BackoffPolicyArg;
22use google_cloud_gax::retry_policy::RetryPolicyArg;
23use std::time::Duration;
24
25/// A builder for [Write].
26///
27/// # Example
28/// ```
29/// # use google_cloud_bigquery::client::Write;
30/// # async fn sample() -> anyhow::Result<()> {
31/// let builder = Write::builder();
32/// let client = builder
33///     .with_endpoint("https://bigquerystoragewrite.googleapis.com")
34///     .build()
35///     .await?;
36/// # Ok(()) }
37/// ```
38#[derive(Debug)]
39pub struct ClientBuilder {
40    pub(super) config: ClientConfig,
41    pub(super) pool_options: StreamPoolOptions,
42    pub(super) retry_options: RetryOptions,
43}
44
45impl ClientBuilder {
46    pub(super) fn new() -> Self {
47        Self {
48            config: ClientConfig::default(),
49            pool_options: StreamPoolOptions::default(),
50            retry_options: RetryOptions::default(),
51        }
52    }
53
54    /// Creates a new client.
55    ///
56    /// # Example
57    /// ```
58    /// # use google_cloud_bigquery::client::Write;
59    /// # async fn sample() -> anyhow::Result<()> {
60    /// let client = Write::builder().build().await?;
61    /// # Ok(()) }
62    /// ```
63    pub async fn build(self) -> BuilderResult<Write> {
64        Write::new(self).await
65    }
66
67    /// Sets the endpoint.
68    ///
69    /// # Example
70    /// ```
71    /// # use google_cloud_bigquery::client::Write;
72    /// # async fn sample() -> anyhow::Result<()> {
73    /// let client = Write::builder()
74    ///     .with_endpoint("https://private.googleapis.com")
75    ///     .build()
76    ///     .await?;
77    /// # Ok(()) }
78    /// ```
79    pub fn with_endpoint<V: Into<String>>(mut self, v: V) -> Self {
80        self.config.endpoint = Some(v.into());
81        self
82    }
83
84    /// Configure the universe domain.
85    ///
86    /// The universe domain is the default service domain for a given cloud universe.
87    /// The default value is "googleapis.com".
88    ///
89    /// # Example
90    /// ```
91    /// # use google_cloud_bigquery::client::Write;
92    /// # async fn sample() -> anyhow::Result<()> {
93    /// let client = Write::builder()
94    ///     .with_universe_domain("googleapis.com")
95    ///     .build()
96    ///     .await?;
97    /// # Ok(()) }
98    /// ```
99    pub fn with_universe_domain<V: Into<String>>(mut self, v: V) -> Self {
100        self.config.universe_domain = Some(v.into());
101        self
102    }
103
104    /// Configures the authentication credentials.
105    ///
106    /// More information about valid credentials types can be found in the
107    /// [google-cloud-auth] crate documentation.
108    ///
109    /// # Example
110    /// ```
111    /// # use google_cloud_bigquery::client::Write;
112    /// # async fn sample() -> anyhow::Result<()> {
113    /// use google_cloud_auth::credentials::mds;
114    /// let client = Write::builder()
115    ///     .with_credentials(
116    ///         mds::Builder::default()
117    ///             .with_scopes(["https://www.googleapis.com/auth/cloud-platform.read-only"])
118    ///             .build()?)
119    ///     .build()
120    ///     .await?;
121    /// # Ok(()) }
122    /// ```
123    ///
124    /// [google-cloud-auth]: https://docs.rs/google-cloud-auth
125    pub fn with_credentials<V: Into<Credentials>>(mut self, v: V) -> Self {
126        self.config.cred = Some(v.into());
127        self
128    }
129
130    /// Configure the number of subchannels used by the client.
131    ///
132    /// # Example
133    /// ```
134    /// # use google_cloud_bigquery::client::Write;
135    /// # async fn sample() -> anyhow::Result<()> {
136    /// let count = std::thread::available_parallelism()?.get();
137    /// let client = Write::builder()
138    ///     .with_grpc_subchannel_count(count)
139    ///     .build()
140    ///     .await?;
141    /// # Ok(()) }
142    /// ```
143    ///
144    /// gRPC-based clients may exhibit high latency if many requests need to be
145    /// demuxed over a single HTTP/2 connection (often called a *subchannel* in
146    /// gRPC).
147    ///
148    /// Consider using more subchannels if your application creates many
149    /// writers. Consider using fewer subchannels if your application needs the
150    /// file descriptors for other purposes.
151    pub fn with_grpc_subchannel_count(mut self, v: usize) -> Self {
152        self.config.grpc_subchannel_count = Some(v);
153        self
154    }
155
156    /// Configure the maximum streams in the client's multiplexed stream pool.
157    ///
158    /// # Example
159    /// ```
160    /// # use google_cloud_bigquery::client::Write;
161    /// # async fn sample() -> anyhow::Result<()> {
162    /// let count = std::thread::available_parallelism()?.get();
163    /// let client = Write::builder()
164    ///     .with_pool_size_limit(count)
165    ///     .build()
166    ///     .await?;
167    /// # Ok(()) }
168    /// ```
169    ///
170    /// This stream pool is shared by default writers with multiplexing enabled.
171    ///
172    /// The client scales the stream pool up to this limit as the streams in the
173    /// pool encounter load.
174    ///
175    /// The default is 8 streams.
176    pub fn with_pool_size_limit(mut self, v: usize) -> Self {
177        self.pool_options.max_streams = v.max(1);
178        self
179    }
180
181    // TODO(#6866) - expose when we have rebalancing
182    #[cfg_attr(not(test), expect(dead_code))]
183    /// Configure the maximum outstanding requests in the client's multiplexed
184    /// stream pool.
185    ///
186    /// # Example
187    /// ```no_rust
188    /// # use google_cloud_bigquery::client::Write;
189    /// # async fn sample() -> anyhow::Result<()> {
190    /// let client = Write::builder()
191    ///     .with_max_outstanding_requests(200)
192    ///     .build()
193    ///     .await?;
194    /// # Ok(()) }
195    /// ```
196    ///
197    /// As streams in the stream pool approach this limit, the client
198    /// dynamically adds more streams to the stream pool, up to the limit
199    /// configured by `with_pool_size_limit`.
200    ///
201    /// The default is 1000 requests.
202    pub(crate) fn with_max_outstanding_requests(mut self, v: u64) -> Self {
203        self.pool_options.max_outstanding_requests = Some(v.max(1));
204        self
205    }
206
207    // TODO(#6866) - expose when we have rebalancing
208    #[cfg_attr(not(test), expect(dead_code))]
209    /// Configure the maximum outstanding bytes in the client's multiplexed
210    /// stream pool.
211    ///
212    /// # Example
213    /// ```no_rust
214    /// # use google_cloud_bigquery::client::Write;
215    /// # async fn sample() -> anyhow::Result<()> {
216    /// let client = Write::builder()
217    ///     .with_max_outstanding_bytes(200_000)
218    ///     .build()
219    ///     .await?;
220    /// # Ok(()) }
221    /// ```
222    ///
223    /// As streams in the stream pool approach this limit, the client
224    /// dynamically adds more streams to the stream pool, up to the limit
225    /// configured by `with_pool_size_limit`.
226    pub(crate) fn with_max_outstanding_bytes(mut self, v: u64) -> Self {
227        self.pool_options.max_outstanding_bytes = Some(v.max(1));
228        self
229    }
230
231    // TODO(#6851) - release when all streams support retries.
232    #[cfg_attr(not(test), expect(dead_code))]
233    /// Configure the retry policy.
234    ///
235    /// The client libraries can automatically retry operations that fail. The
236    /// retry policy controls what errors are considered retryable, sets limits
237    /// on the number of attempts or the time trying to make attempts.
238    ///
239    /// # Example
240    /// ```no_rust
241    /// # use google_cloud_bigquery::client::Write;
242    /// # async fn sample() -> anyhow::Result<()> {
243    /// use google_cloud_bigquery::write::retry_policy::RetryableErrors;
244    /// use google_cloud_gax::retry_policy::RetryPolicyExt;
245    /// let client = Write::builder()
246    ///     .with_retry_policy(RetryableErrors.with_attempt_limit(3))
247    ///     .build()
248    ///     .await?;
249    /// # Ok(()) }
250    /// ```
251    pub(crate) fn with_retry_policy<V: Into<RetryPolicyArg>>(mut self, v: V) -> Self {
252        self.retry_options.retry_policy = v.into().into();
253        self
254    }
255
256    // TODO(#6851) - release when all streams support retries.
257    #[cfg_attr(not(test), expect(dead_code))]
258    /// Configure the retry backoff policy.
259    ///
260    /// The client libraries can automatically retry operations that fail. The
261    /// backoff policy controls how long to wait in between retry attempts.
262    ///
263    /// # Example
264    /// ```no_rust
265    /// # use google_cloud_bigquery::client::Write;
266    /// # async fn sample() -> anyhow::Result<()> {
267    /// use google_cloud_gax::exponential_backoff::ExponentialBackoff;
268    /// let policy = ExponentialBackoff::default();
269    /// let client = Write::builder()
270    ///     .with_backoff_policy(policy)
271    ///     .build()
272    ///     .await?;
273    /// # Ok(()) }
274    /// ```
275    pub(crate) fn with_backoff_policy<V: Into<BackoffPolicyArg>>(mut self, v: V) -> Self {
276        self.retry_options.backoff_policy = v.into().into();
277        self
278    }
279
280    // TODO(#6851) - release when all streams support retries.
281    #[cfg_attr(not(test), expect(dead_code))]
282    /// Configure the timeout for a single write attempt.
283    ///
284    /// Without this limit, a write can block forever if the service accepts
285    /// the stream but never responds. On a timeout, the client abandons the
286    /// stream and the retry policy decides whether to make another attempt.
287    ///
288    /// # Example
289    /// ```no_rust
290    /// # use google_cloud_bigquery::client::Write;
291    /// # async fn sample() -> anyhow::Result<()> {
292    /// use std::time::Duration;
293    /// let client = Write::builder()
294    ///     .with_attempt_timeout(Duration::from_secs(10))
295    ///     .build()
296    ///     .await?;
297    /// # Ok(()) }
298    /// ```
299    pub(crate) fn with_attempt_timeout(mut self, v: Duration) -> Self {
300        self.retry_options.attempt_timeout = Some(v);
301        self
302    }
303}
304
305#[cfg(test)]
306mod tests {
307    use super::*;
308    use crate::write::test::NoBackoff;
309    use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
310    use google_cloud_gax::retry_policy::NeverRetry;
311
312    #[test]
313    fn defaults() {
314        let builder = ClientBuilder::new();
315        assert!(builder.config.endpoint.is_none(), "{:?}", builder.config);
316        assert!(builder.config.cred.is_none(), "{:?}", builder.config);
317        assert!(
318            builder.config.universe_domain.is_none(),
319            "{:?}",
320            builder.config
321        );
322        assert!(
323            builder.config.grpc_subchannel_count.is_none(),
324            "{:?}",
325            builder.config
326        );
327        assert_eq!(builder.pool_options.max_streams, 8);
328        assert_eq!(builder.pool_options.max_outstanding_requests, Some(1000));
329        assert_eq!(builder.pool_options.max_outstanding_bytes, None);
330    }
331
332    #[test]
333    fn setters() {
334        let builder = ClientBuilder::new()
335            .with_endpoint("test-endpoint.com")
336            .with_universe_domain("test-ud.com")
337            .with_credentials(Anonymous::new().build())
338            .with_grpc_subchannel_count(16)
339            .with_pool_size_limit(10)
340            .with_max_outstanding_requests(900)
341            .with_max_outstanding_bytes(1_000_000)
342            .with_retry_policy(NeverRetry)
343            .with_backoff_policy(NoBackoff)
344            .with_attempt_timeout(Duration::from_secs(10));
345        assert_eq!(
346            builder.config.endpoint,
347            Some("test-endpoint.com".to_string())
348        );
349        assert_eq!(
350            builder.config.universe_domain,
351            Some("test-ud.com".to_string())
352        );
353        assert!(builder.config.cred.is_some(), "{:?}", builder.config);
354        assert_eq!(builder.config.grpc_subchannel_count, Some(16));
355        assert_eq!(builder.pool_options.max_streams, 10);
356        assert_eq!(builder.pool_options.max_outstanding_requests, Some(900));
357        assert_eq!(builder.pool_options.max_outstanding_bytes, Some(1_000_000));
358        assert_eq!(
359            builder.retry_options.attempt_timeout,
360            Some(Duration::from_secs(10))
361        );
362
363        let fmt = format!("{:?}", builder.retry_options);
364        assert!(fmt.contains("NeverRetry"), "{fmt}");
365        assert!(fmt.contains("NoBackoff"), "{fmt}");
366    }
367
368    #[test]
369    fn validate_pool_options() {
370        let builder = ClientBuilder::new()
371            .with_credentials(Anonymous::new().build())
372            .with_pool_size_limit(0)
373            .with_max_outstanding_requests(0)
374            .with_max_outstanding_bytes(0);
375        assert_eq!(builder.pool_options.max_streams, 1);
376        assert_eq!(builder.pool_options.max_outstanding_requests, Some(1));
377        assert_eq!(builder.pool_options.max_outstanding_bytes, Some(1));
378    }
379}