google_cloud_bigquery/write/generated/gapic_storage/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//
15// Code generated by sidekick. DO NOT EDIT.
16#![allow(rustdoc::bare_urls)]
17#![allow(rustdoc::broken_intra_doc_links)]
18#![allow(rustdoc::invalid_html_tags)]
19#![allow(rustdoc::redundant_explicit_links)]
20
21/// Implements a client for the BigQuery Storage API.
22///
23/// # Example
24/// ```
25/// # use google_cloud_bigquery::client::Read;
26/// async fn sample(
27/// ) -> anyhow::Result<()> {
28/// let client = Read::builder().build().await?;
29/// let response = client.create_read_session()
30/// /* set fields */
31/// .send().await?;
32/// println!("response {:?}", response);
33/// Ok(())
34/// }
35/// ```
36///
37/// # Service Description
38///
39/// BigQuery Read API.
40///
41/// The Read API can be used to read data from BigQuery.
42///
43/// # Configuration
44///
45/// To configure `Read` use the `with_*` methods in the type returned
46/// by [builder()][Read::builder]. The default configuration should
47/// work for most applications. Common configuration changes include
48///
49/// * [with_endpoint()]: by default this client uses the global default endpoint
50/// (`https://bigquerystorage.googleapis.com`). Applications using regional
51/// endpoints or running in restricted networks (e.g. a network configured
52/// with [Private Google Access with VPC Service Controls]) may want to
53/// override this default.
54/// * [with_credentials()]: by default this client uses
55/// [Application Default Credentials]. Applications using custom
56/// authentication may need to override this default.
57///
58/// [with_endpoint()]: super::builder::read::ClientBuilder::with_endpoint
59/// [with_credentials()]: super::builder::read::ClientBuilder::with_credentials
60/// [Private Google Access with VPC Service Controls]: https://cloud.google.com/vpc-service-controls/docs/private-connectivity
61/// [Application Default Credentials]: https://cloud.google.com/docs/authentication#adc
62///
63/// # Pooling and Cloning
64///
65/// `Read` holds a connection pool internally, it is advised to
66/// create one and reuse it. You do not need to wrap `Read` in
67/// an [Rc](std::rc::Rc) or [Arc](std::sync::Arc) to reuse it, because it
68/// already uses an `Arc` internally.
69#[derive(Clone, Debug)]
70pub struct Read {
71 inner: std::sync::Arc<dyn super::stub::dynamic::Read>,
72}
73
74impl Read {
75 /// Returns a builder for [Read].
76 ///
77 /// ```
78 /// # async fn sample() -> google_cloud_gax::client_builder::Result<()> {
79 /// # use google_cloud_bigquery::client::Read;
80 /// let client = Read::builder().build().await?;
81 /// # Ok(()) }
82 /// ```
83 pub fn builder() -> super::builder::read::ClientBuilder {
84 crate::new_client_builder(super::builder::read::client::Factory)
85 }
86
87 /// Creates a new client from the provided stub.
88 ///
89 /// The most common case for calling this function is in tests mocking the
90 /// client's behavior.
91 pub fn from_stub<T>(stub: impl Into<std::sync::Arc<T>>) -> Self
92 where
93 T: super::stub::Read + 'static,
94 {
95 Self { inner: stub.into() }
96 }
97
98 pub(crate) async fn new(
99 config: gaxi::options::ClientConfig,
100 ) -> crate::ClientBuilderResult<Self> {
101 let inner = Self::build_inner(config).await?;
102 Ok(Self { inner })
103 }
104
105 async fn build_inner(
106 conf: gaxi::options::ClientConfig,
107 ) -> crate::ClientBuilderResult<std::sync::Arc<dyn super::stub::dynamic::Read>> {
108 if gaxi::options::tracing_enabled(&conf) {
109 return Ok(std::sync::Arc::new(Self::build_with_tracing(conf).await?));
110 }
111 Ok(std::sync::Arc::new(Self::build_transport(conf).await?))
112 }
113
114 async fn build_transport(
115 conf: gaxi::options::ClientConfig,
116 ) -> crate::ClientBuilderResult<impl super::stub::Read> {
117 super::transport::Read::new(conf).await
118 }
119
120 async fn build_with_tracing(
121 conf: gaxi::options::ClientConfig,
122 ) -> crate::ClientBuilderResult<impl super::stub::Read> {
123 Self::build_transport(conf)
124 .await
125 .map(super::tracing::Read::new)
126 }
127
128 /// Creates a new read session. A read session divides the contents of a
129 /// BigQuery table into one or more streams, which can then be used to read
130 /// data from the table. The read session also specifies properties of the
131 /// data to be read, such as a list of columns or a push-down filter describing
132 /// the rows to be returned.
133 ///
134 /// A particular row can be read by at most one stream. When the caller has
135 /// reached the end of each stream in the session, then all the data in the
136 /// table has been read.
137 ///
138 /// Data is assigned to each stream such that roughly the same number of
139 /// rows can be read from each stream. Because the server-side unit for
140 /// assigning data is collections of rows, the API does not guarantee that
141 /// each stream will return the same number or rows. Additionally, the
142 /// limits are enforced based on the number of pre-filtered rows, so some
143 /// filters can lead to lopsided assignments.
144 ///
145 /// Read sessions automatically expire 6 hours after they are created and do
146 /// not require manual clean-up by the caller.
147 ///
148 /// # Example
149 /// ```
150 /// # use google_cloud_bigquery::client::Read;
151 /// use google_cloud_bigquery::Result;
152 /// async fn sample(
153 /// client: &Read
154 /// ) -> Result<()> {
155 /// let response = client.create_read_session()
156 /// /* set fields */
157 /// .send().await?;
158 /// println!("response {:?}", response);
159 /// Ok(())
160 /// }
161 /// ```
162 pub fn create_read_session(&self) -> super::builder::read::CreateReadSession {
163 super::builder::read::CreateReadSession::new(self.inner.clone())
164 }
165
166 /// Reads rows from the stream in the format prescribed by the ReadSession.
167 /// Each response contains one or more table rows, up to a maximum of 128 MB
168 /// per response; read requests which attempt to read individual rows larger
169 /// than 128 MB will fail.
170 ///
171 /// Each request also returns a set of stream statistics reflecting the current
172 /// state of the stream.
173 ///
174 /// # Example
175 /// ```
176 /// # use google_cloud_bigquery::client::Read;
177 /// use google_cloud_bigquery::Result;
178 /// async fn sample(
179 /// client: &Read
180 /// ) -> Result<()> {
181 /// let mut resp_stream = client.read_rows()
182 /// /* set fields */
183 /// .send().await?;
184 /// while let Some(response) = resp_stream.next().await {
185 /// let response = response?;
186 /// println!("response {:?}", response);
187 /// }
188 /// Ok(())
189 /// }
190 /// ```
191 pub fn read_rows(&self) -> super::builder::read::ReadRows {
192 super::builder::read::ReadRows::new(self.inner.clone())
193 }
194
195 /// Splits a given `ReadStream` into two `ReadStream` objects. These
196 /// `ReadStream` objects are referred to as the primary and the residual
197 /// streams of the split. The original `ReadStream` can still be read from in
198 /// the same manner as before. Both of the returned `ReadStream` objects can
199 /// also be read from, and the rows returned by both child streams will be
200 /// the same as the rows read from the original stream.
201 ///
202 /// Moreover, the two child streams will be allocated back-to-back in the
203 /// original `ReadStream`. Concretely, it is guaranteed that for streams
204 /// original, primary, and residual, that original[0-j] = primary[0-j] and
205 /// original[j-n] = residual[0-m] once the streams have been read to
206 /// completion.
207 ///
208 /// # Example
209 /// ```
210 /// # use google_cloud_bigquery::client::Read;
211 /// use google_cloud_bigquery::Result;
212 /// async fn sample(
213 /// client: &Read
214 /// ) -> Result<()> {
215 /// let response = client.split_read_stream()
216 /// /* set fields */
217 /// .send().await?;
218 /// println!("response {:?}", response);
219 /// Ok(())
220 /// }
221 /// ```
222 pub fn split_read_stream(&self) -> super::builder::read::SplitReadStream {
223 super::builder::read::SplitReadStream::new(self.inner.clone())
224 }
225}
226
227/// Implements a client for the BigQuery Storage API.
228#[derive(Clone, Debug)]
229pub struct BigQueryWrite {
230 inner: std::sync::Arc<dyn super::stub::dynamic::BigQueryWrite>,
231}
232
233impl BigQueryWrite {
234 /// Creates a new client from the provided stub.
235 ///
236 /// The most common case for calling this function is in tests mocking the
237 /// client's behavior.
238 pub fn from_stub<T>(stub: impl Into<std::sync::Arc<T>>) -> Self
239 where
240 T: super::stub::BigQueryWrite + 'static,
241 {
242 Self { inner: stub.into() }
243 }
244
245 pub(crate) async fn new(
246 config: gaxi::options::ClientConfig,
247 ) -> crate::ClientBuilderResult<Self> {
248 let inner = Self::build_inner(config).await?;
249 Ok(Self { inner })
250 }
251
252 async fn build_inner(
253 conf: gaxi::options::ClientConfig,
254 ) -> crate::ClientBuilderResult<std::sync::Arc<dyn super::stub::dynamic::BigQueryWrite>> {
255 if gaxi::options::tracing_enabled(&conf) {
256 return Ok(std::sync::Arc::new(Self::build_with_tracing(conf).await?));
257 }
258 Ok(std::sync::Arc::new(Self::build_transport(conf).await?))
259 }
260
261 async fn build_transport(
262 conf: gaxi::options::ClientConfig,
263 ) -> crate::ClientBuilderResult<impl super::stub::BigQueryWrite> {
264 super::transport::BigQueryWrite::new(conf).await
265 }
266
267 async fn build_with_tracing(
268 conf: gaxi::options::ClientConfig,
269 ) -> crate::ClientBuilderResult<impl super::stub::BigQueryWrite> {
270 Self::build_transport(conf)
271 .await
272 .map(super::tracing::BigQueryWrite::new)
273 }
274
275 /// Creates a write stream to the given table.
276 /// Additionally, every table has a special stream named '_default'
277 /// to which data can be written. This stream doesn't need to be created using
278 /// CreateWriteStream. It is a stream that can be used simultaneously by any
279 /// number of clients. Data written to this stream is considered committed as
280 /// soon as an acknowledgement is received.
281 pub(crate) fn create_write_stream(&self) -> super::builder::big_query_write::CreateWriteStream {
282 super::builder::big_query_write::CreateWriteStream::new(self.inner.clone())
283 }
284
285 /// Appends data to the given stream.
286 ///
287 /// If `offset` is specified, the `offset` is checked against the end of
288 /// stream. The server returns `OUT_OF_RANGE` in `AppendRowsResponse` if an
289 /// attempt is made to append to an offset beyond the current end of the stream
290 /// or `ALREADY_EXISTS` if user provides an `offset` that has already been
291 /// written to. User can retry with adjusted offset within the same RPC
292 /// connection. If `offset` is not specified, append happens at the end of the
293 /// stream.
294 ///
295 /// The response contains an optional offset at which the append
296 /// happened. No offset information will be returned for appends to a
297 /// default stream.
298 ///
299 /// Responses are received in the same order in which requests are sent.
300 /// There will be one response for each successful inserted request. Responses
301 /// may optionally embed error information if the originating AppendRequest was
302 /// not successfully processed.
303 ///
304 /// The specifics of when successfully appended data is made visible to the
305 /// table are governed by the type of stream:
306 ///
307 /// * For COMMITTED streams (which includes the default stream), data is
308 /// visible immediately upon successful append.
309 ///
310 /// * For BUFFERED streams, data is made visible via a subsequent `FlushRows`
311 /// rpc which advances a cursor to a newer offset in the stream.
312 ///
313 /// * For PENDING streams, data is not made visible until the stream itself is
314 /// finalized (via the `FinalizeWriteStream` rpc), and the stream is explicitly
315 /// committed via the `BatchCommitWriteStreams` rpc.
316 ///
317 pub(crate) fn append_rows(&self) -> super::builder::big_query_write::AppendRows {
318 super::builder::big_query_write::AppendRows::new(self.inner.clone())
319 }
320
321 /// Gets information about a write stream.
322 pub(crate) fn get_write_stream(&self) -> super::builder::big_query_write::GetWriteStream {
323 super::builder::big_query_write::GetWriteStream::new(self.inner.clone())
324 }
325
326 /// Finalize a write stream so that no new data can be appended to the
327 /// stream. Finalize is not supported on the '_default' stream.
328 pub(crate) fn finalize_write_stream(
329 &self,
330 ) -> super::builder::big_query_write::FinalizeWriteStream {
331 super::builder::big_query_write::FinalizeWriteStream::new(self.inner.clone())
332 }
333
334 /// Atomically commits a group of `PENDING` streams that belong to the same
335 /// `parent` table.
336 ///
337 /// Streams must be finalized before commit and cannot be committed multiple
338 /// times. Once a stream is committed, data in the stream becomes available
339 /// for read operations.
340 pub(crate) fn batch_commit_write_streams(
341 &self,
342 ) -> super::builder::big_query_write::BatchCommitWriteStreams {
343 super::builder::big_query_write::BatchCommitWriteStreams::new(self.inner.clone())
344 }
345
346 /// Flushes rows to a BUFFERED stream.
347 ///
348 /// If users are appending rows to BUFFERED stream, flush operation is
349 /// required in order for the rows to become available for reading. A
350 /// Flush operation flushes up to any previously flushed offset in a BUFFERED
351 /// stream, to the offset specified in the request.
352 ///
353 /// Flush is not supported on the _default stream, since it is not BUFFERED.
354 pub(crate) fn flush_rows(&self) -> super::builder::big_query_write::FlushRows {
355 super::builder::big_query_write::FlushRows::new(self.inner.clone())
356 }
357}