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://bigquerystorage.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/bigquery"])
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    /// Note that until background stream rebalancing ([#6866]) is implemented,
176    /// stream assignment and pool scale-up only occur when a writer is built.
177    /// If multiple writers are created all at once before writes are in flight,
178    /// they will all share the same initial stream. To scale across multiple
179    /// streams, stagger writer creation so writes are in flight when new
180    /// writers are built.
181    ///
182    /// The default is 8 streams.
183    ///
184    /// [#6866]: https://github.com/googleapis/google-cloud-rust/issues/6866
185    pub fn with_pool_size_limit(mut self, v: usize) -> Self {
186        self.pool_options.max_streams = v.max(1);
187        self
188    }
189
190    // TODO(#6866) - expose when we have rebalancing
191    #[cfg_attr(not(test), expect(dead_code))]
192    /// Configure the maximum outstanding requests per stream in the client's
193    /// multiplexed stream pool.
194    ///
195    /// # Example
196    /// ```no_rust
197    /// # use google_cloud_bigquery::client::Write;
198    /// # async fn sample() -> anyhow::Result<()> {
199    /// let client = Write::builder()
200    ///     .with_max_outstanding_requests(200)
201    ///     .build()
202    ///     .await?;
203    /// # Ok(()) }
204    /// ```
205    ///
206    /// As streams in the stream pool approach this limit, the client
207    /// dynamically adds more streams to the stream pool, up to the limit
208    /// configured by `with_pool_size_limit`.
209    ///
210    /// Note that this is not a hard maximum. When the pool reaches the stream
211    /// limit, the streams will continue accepting requests. Consider using a
212    /// semaphore locally for flow control.
213    ///
214    /// The default is 1 request.
215    pub(crate) fn with_max_outstanding_requests(mut self, v: u64) -> Self {
216        self.pool_options.max_outstanding_requests = Some(v.max(1));
217        self
218    }
219
220    // TODO(#6866) - expose when we have rebalancing
221    #[cfg_attr(not(test), expect(dead_code))]
222    /// Configure the maximum outstanding bytes per stream in the client's
223    /// multiplexed stream pool.
224    ///
225    /// # Example
226    /// ```no_rust
227    /// # use google_cloud_bigquery::client::Write;
228    /// # async fn sample() -> anyhow::Result<()> {
229    /// let client = Write::builder()
230    ///     .with_max_outstanding_bytes(200_000)
231    ///     .build()
232    ///     .await?;
233    /// # Ok(()) }
234    /// ```
235    ///
236    /// As streams in the stream pool approach this limit, the client
237    /// dynamically adds more streams to the stream pool, up to the limit
238    /// configured by `with_pool_size_limit`.
239    ///
240    /// Note that this is not a hard maximum. When the pool reaches the stream
241    /// limit, the streams will continue accepting requests. Consider using a
242    /// semaphore locally for flow control.
243    pub(crate) fn with_max_outstanding_bytes(mut self, v: u64) -> Self {
244        self.pool_options.max_outstanding_bytes = Some(v.max(1));
245        self
246    }
247
248    // TODO(#6851) - release when all streams support retries.
249    #[cfg_attr(not(test), expect(dead_code))]
250    /// Configure the retry policy.
251    ///
252    /// The client libraries can automatically retry operations that fail. The
253    /// retry policy controls what errors are considered retryable, sets limits
254    /// on the number of attempts or the time trying to make attempts.
255    ///
256    /// # Example
257    /// ```no_rust
258    /// # use google_cloud_bigquery::client::Write;
259    /// # async fn sample() -> anyhow::Result<()> {
260    /// use google_cloud_bigquery::write::retry_policy::RetryableErrors;
261    /// use google_cloud_gax::retry_policy::RetryPolicyExt;
262    /// let client = Write::builder()
263    ///     .with_retry_policy(RetryableErrors.with_attempt_limit(3))
264    ///     .build()
265    ///     .await?;
266    /// # Ok(()) }
267    /// ```
268    pub(crate) fn with_retry_policy<V: Into<RetryPolicyArg>>(mut self, v: V) -> Self {
269        self.retry_options.retry_policy = v.into().into();
270        self
271    }
272
273    // TODO(#6851) - release when all streams support retries.
274    #[cfg_attr(not(test), expect(dead_code))]
275    /// Configure the retry backoff policy.
276    ///
277    /// The client libraries can automatically retry operations that fail. The
278    /// backoff policy controls how long to wait in between retry attempts.
279    ///
280    /// # Example
281    /// ```no_rust
282    /// # use google_cloud_bigquery::client::Write;
283    /// # async fn sample() -> anyhow::Result<()> {
284    /// use google_cloud_gax::exponential_backoff::ExponentialBackoff;
285    /// let policy = ExponentialBackoff::default();
286    /// let client = Write::builder()
287    ///     .with_backoff_policy(policy)
288    ///     .build()
289    ///     .await?;
290    /// # Ok(()) }
291    /// ```
292    pub(crate) fn with_backoff_policy<V: Into<BackoffPolicyArg>>(mut self, v: V) -> Self {
293        self.retry_options.backoff_policy = v.into().into();
294        self
295    }
296
297    // TODO(#6851) - release when all streams support retries.
298    #[cfg_attr(not(test), expect(dead_code))]
299    /// Configure the timeout for a single write attempt.
300    ///
301    /// Without this limit, a write can block forever if the service accepts
302    /// the stream but never responds. On a timeout, the client abandons the
303    /// stream and the retry policy decides whether to make another attempt.
304    ///
305    /// # Example
306    /// ```no_rust
307    /// # use google_cloud_bigquery::client::Write;
308    /// # async fn sample() -> anyhow::Result<()> {
309    /// use std::time::Duration;
310    /// let client = Write::builder()
311    ///     .with_attempt_timeout(Duration::from_secs(10))
312    ///     .build()
313    ///     .await?;
314    /// # Ok(()) }
315    /// ```
316    pub(crate) fn with_attempt_timeout(mut self, v: Duration) -> Self {
317        self.retry_options.attempt_timeout = Some(v);
318        self
319    }
320}
321
322#[cfg(test)]
323mod tests {
324    use super::*;
325    use crate::write::test::NoBackoff;
326    use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
327    use google_cloud_gax::retry_policy::NeverRetry;
328
329    #[test]
330    fn defaults() {
331        let builder = ClientBuilder::new();
332        assert!(builder.config.endpoint.is_none(), "{:?}", builder.config);
333        assert!(builder.config.cred.is_none(), "{:?}", builder.config);
334        assert!(
335            builder.config.universe_domain.is_none(),
336            "{:?}",
337            builder.config
338        );
339        assert!(
340            builder.config.grpc_subchannel_count.is_none(),
341            "{:?}",
342            builder.config
343        );
344        assert_eq!(builder.pool_options.max_streams, 8);
345        assert_eq!(builder.pool_options.max_outstanding_requests, Some(1));
346        assert_eq!(builder.pool_options.max_outstanding_bytes, None);
347    }
348
349    #[test]
350    fn setters() {
351        let builder = ClientBuilder::new()
352            .with_endpoint("test-endpoint.com")
353            .with_universe_domain("test-ud.com")
354            .with_credentials(Anonymous::new().build())
355            .with_grpc_subchannel_count(16)
356            .with_pool_size_limit(10)
357            .with_max_outstanding_requests(900)
358            .with_max_outstanding_bytes(1_000_000)
359            .with_retry_policy(NeverRetry)
360            .with_backoff_policy(NoBackoff)
361            .with_attempt_timeout(Duration::from_secs(10));
362        assert_eq!(
363            builder.config.endpoint,
364            Some("test-endpoint.com".to_string())
365        );
366        assert_eq!(
367            builder.config.universe_domain,
368            Some("test-ud.com".to_string())
369        );
370        assert!(builder.config.cred.is_some(), "{:?}", builder.config);
371        assert_eq!(builder.config.grpc_subchannel_count, Some(16));
372        assert_eq!(builder.pool_options.max_streams, 10);
373        assert_eq!(builder.pool_options.max_outstanding_requests, Some(900));
374        assert_eq!(builder.pool_options.max_outstanding_bytes, Some(1_000_000));
375        assert_eq!(
376            builder.retry_options.attempt_timeout,
377            Some(Duration::from_secs(10))
378        );
379
380        let fmt = format!("{:?}", builder.retry_options);
381        assert!(fmt.contains("NeverRetry"), "{fmt}");
382        assert!(fmt.contains("NoBackoff"), "{fmt}");
383    }
384
385    #[test]
386    fn validate_pool_options() {
387        let builder = ClientBuilder::new()
388            .with_credentials(Anonymous::new().build())
389            .with_pool_size_limit(0)
390            .with_max_outstanding_requests(0)
391            .with_max_outstanding_bytes(0);
392        assert_eq!(builder.pool_options.max_streams, 1);
393        assert_eq!(builder.pool_options.max_outstanding_requests, Some(1));
394        assert_eq!(builder.pool_options.max_outstanding_bytes, Some(1));
395    }
396}