Skip to main content

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}