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}