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}