google-cloud-bigquery 0.17.0-preview

Google Cloud Client Libraries for Rust - BigQuery
Documentation
// Copyright 2026 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
//     https://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//
// Code generated by sidekick. DO NOT EDIT.
#![allow(rustdoc::bare_urls)]
#![allow(rustdoc::broken_intra_doc_links)]
#![allow(rustdoc::invalid_html_tags)]
#![allow(rustdoc::redundant_explicit_links)]

/// Implements a client for the BigQuery Storage API.
///
/// # Example
/// ```
/// # use google_cloud_bigquery::client::Read;
/// async fn sample(
/// ) -> anyhow::Result<()> {
///     let client = Read::builder().build().await?;
///     let response = client.create_read_session()
///         /* set fields */
///         .send().await?;
///     println!("response {:?}", response);
///     Ok(())
/// }
/// ```
///
/// # Service Description
///
/// BigQuery Read API.
///
/// The Read API can be used to read data from BigQuery.
///
/// # Configuration
///
/// To configure `Read` use the `with_*` methods in the type returned
/// by [builder()][Read::builder]. The default configuration should
/// work for most applications. Common configuration changes include
///
/// * [with_endpoint()]: by default this client uses the global default endpoint
///   (`https://bigquerystorage.googleapis.com`). Applications using regional
///   endpoints or running in restricted networks (e.g. a network configured
///   with [Private Google Access with VPC Service Controls]) may want to
///   override this default.
/// * [with_credentials()]: by default this client uses
///   [Application Default Credentials]. Applications using custom
///   authentication may need to override this default.
///
/// [with_endpoint()]: super::builder::read::ClientBuilder::with_endpoint
/// [with_credentials()]: super::builder::read::ClientBuilder::with_credentials
/// [Private Google Access with VPC Service Controls]: https://cloud.google.com/vpc-service-controls/docs/private-connectivity
/// [Application Default Credentials]: https://cloud.google.com/docs/authentication#adc
///
/// # Pooling and Cloning
///
/// `Read` holds a connection pool internally, it is advised to
/// create one and reuse it. You do not need to wrap `Read` in
/// an [Rc](std::rc::Rc) or [Arc](std::sync::Arc) to reuse it, because it
/// already uses an `Arc` internally.
#[derive(Clone, Debug)]
pub struct Read {
    inner: std::sync::Arc<dyn super::stub::dynamic::Read>,
}

impl Read {
    /// Returns a builder for [Read].
    ///
    /// ```
    /// # async fn sample() -> google_cloud_gax::client_builder::Result<()> {
    /// # use google_cloud_bigquery::client::Read;
    /// let client = Read::builder().build().await?;
    /// # Ok(()) }
    /// ```
    pub fn builder() -> super::builder::read::ClientBuilder {
        crate::new_client_builder(super::builder::read::client::Factory)
    }

    /// Creates a new client from the provided stub.
    ///
    /// The most common case for calling this function is in tests mocking the
    /// client's behavior.
    pub fn from_stub<T>(stub: impl Into<std::sync::Arc<T>>) -> Self
    where
        T: super::stub::Read + 'static,
    {
        Self { inner: stub.into() }
    }

    pub(crate) async fn new(
        config: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<Self> {
        let inner = Self::build_inner(config).await?;
        Ok(Self { inner })
    }

    async fn build_inner(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<std::sync::Arc<dyn super::stub::dynamic::Read>> {
        if gaxi::options::tracing_enabled(&conf) {
            return Ok(std::sync::Arc::new(Self::build_with_tracing(conf).await?));
        }
        Ok(std::sync::Arc::new(Self::build_transport(conf).await?))
    }

    async fn build_transport(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<impl super::stub::Read> {
        super::transport::Read::new(conf).await
    }

    async fn build_with_tracing(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<impl super::stub::Read> {
        Self::build_transport(conf)
            .await
            .map(super::tracing::Read::new)
    }

    /// Creates a new read session. A read session divides the contents of a
    /// BigQuery table into one or more streams, which can then be used to read
    /// data from the table. The read session also specifies properties of the
    /// data to be read, such as a list of columns or a push-down filter describing
    /// the rows to be returned.
    ///
    /// A particular row can be read by at most one stream. When the caller has
    /// reached the end of each stream in the session, then all the data in the
    /// table has been read.
    ///
    /// Data is assigned to each stream such that roughly the same number of
    /// rows can be read from each stream. Because the server-side unit for
    /// assigning data is collections of rows, the API does not guarantee that
    /// each stream will return the same number or rows. Additionally, the
    /// limits are enforced based on the number of pre-filtered rows, so some
    /// filters can lead to lopsided assignments.
    ///
    /// Read sessions automatically expire 6 hours after they are created and do
    /// not require manual clean-up by the caller.
    ///
    /// # Example
    /// ```
    /// # use google_cloud_bigquery::client::Read;
    /// use google_cloud_bigquery::Result;
    /// async fn sample(
    ///    client: &Read
    /// ) -> Result<()> {
    ///     let response = client.create_read_session()
    ///         /* set fields */
    ///         .send().await?;
    ///     println!("response {:?}", response);
    ///     Ok(())
    /// }
    /// ```
    pub fn create_read_session(&self) -> super::builder::read::CreateReadSession {
        super::builder::read::CreateReadSession::new(self.inner.clone())
    }

    /// Reads rows from the stream in the format prescribed by the ReadSession.
    /// Each response contains one or more table rows, up to a maximum of 128 MB
    /// per response; read requests which attempt to read individual rows larger
    /// than 128 MB will fail.
    ///
    /// Each request also returns a set of stream statistics reflecting the current
    /// state of the stream.
    ///
    /// # Example
    /// ```
    /// # use google_cloud_bigquery::client::Read;
    /// use google_cloud_bigquery::Result;
    /// async fn sample(
    ///    client: &Read
    /// ) -> Result<()> {
    ///     let mut resp_stream = client.read_rows()
    ///         /* set fields */
    ///         .send().await?;
    ///     while let Some(response) = resp_stream.next().await {
    ///         let response = response?;
    ///         println!("response {:?}", response);
    ///     }
    ///     Ok(())
    /// }
    /// ```
    pub fn read_rows(&self) -> super::builder::read::ReadRows {
        super::builder::read::ReadRows::new(self.inner.clone())
    }

    /// Splits a given `ReadStream` into two `ReadStream` objects. These
    /// `ReadStream` objects are referred to as the primary and the residual
    /// streams of the split. The original `ReadStream` can still be read from in
    /// the same manner as before. Both of the returned `ReadStream` objects can
    /// also be read from, and the rows returned by both child streams will be
    /// the same as the rows read from the original stream.
    ///
    /// Moreover, the two child streams will be allocated back-to-back in the
    /// original `ReadStream`. Concretely, it is guaranteed that for streams
    /// original, primary, and residual, that original[0-j] = primary[0-j] and
    /// original[j-n] = residual[0-m] once the streams have been read to
    /// completion.
    ///
    /// # Example
    /// ```
    /// # use google_cloud_bigquery::client::Read;
    /// use google_cloud_bigquery::Result;
    /// async fn sample(
    ///    client: &Read
    /// ) -> Result<()> {
    ///     let response = client.split_read_stream()
    ///         /* set fields */
    ///         .send().await?;
    ///     println!("response {:?}", response);
    ///     Ok(())
    /// }
    /// ```
    pub fn split_read_stream(&self) -> super::builder::read::SplitReadStream {
        super::builder::read::SplitReadStream::new(self.inner.clone())
    }
}

/// Implements a client for the BigQuery Storage API.
#[derive(Clone, Debug)]
pub struct BigQueryWrite {
    inner: std::sync::Arc<dyn super::stub::dynamic::BigQueryWrite>,
}

impl BigQueryWrite {
    /// Creates a new client from the provided stub.
    ///
    /// The most common case for calling this function is in tests mocking the
    /// client's behavior.
    pub fn from_stub<T>(stub: impl Into<std::sync::Arc<T>>) -> Self
    where
        T: super::stub::BigQueryWrite + 'static,
    {
        Self { inner: stub.into() }
    }

    pub(crate) async fn new(
        config: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<Self> {
        let inner = Self::build_inner(config).await?;
        Ok(Self { inner })
    }

    async fn build_inner(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<std::sync::Arc<dyn super::stub::dynamic::BigQueryWrite>> {
        if gaxi::options::tracing_enabled(&conf) {
            return Ok(std::sync::Arc::new(Self::build_with_tracing(conf).await?));
        }
        Ok(std::sync::Arc::new(Self::build_transport(conf).await?))
    }

    async fn build_transport(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<impl super::stub::BigQueryWrite> {
        super::transport::BigQueryWrite::new(conf).await
    }

    async fn build_with_tracing(
        conf: gaxi::options::ClientConfig,
    ) -> crate::ClientBuilderResult<impl super::stub::BigQueryWrite> {
        Self::build_transport(conf)
            .await
            .map(super::tracing::BigQueryWrite::new)
    }

    /// Creates a write stream to the given table.
    /// Additionally, every table has a special stream named '_default'
    /// to which data can be written. This stream doesn't need to be created using
    /// CreateWriteStream. It is a stream that can be used simultaneously by any
    /// number of clients. Data written to this stream is considered committed as
    /// soon as an acknowledgement is received.
    pub(crate) fn create_write_stream(&self) -> super::builder::big_query_write::CreateWriteStream {
        super::builder::big_query_write::CreateWriteStream::new(self.inner.clone())
    }

    /// Appends data to the given stream.
    ///
    /// If `offset` is specified, the `offset` is checked against the end of
    /// stream. The server returns `OUT_OF_RANGE` in `AppendRowsResponse` if an
    /// attempt is made to append to an offset beyond the current end of the stream
    /// or `ALREADY_EXISTS` if user provides an `offset` that has already been
    /// written to. User can retry with adjusted offset within the same RPC
    /// connection. If `offset` is not specified, append happens at the end of the
    /// stream.
    ///
    /// The response contains an optional offset at which the append
    /// happened.  No offset information will be returned for appends to a
    /// default stream.
    ///
    /// Responses are received in the same order in which requests are sent.
    /// There will be one response for each successful inserted request.  Responses
    /// may optionally embed error information if the originating AppendRequest was
    /// not successfully processed.
    ///
    /// The specifics of when successfully appended data is made visible to the
    /// table are governed by the type of stream:
    ///
    /// * For COMMITTED streams (which includes the default stream), data is
    ///   visible immediately upon successful append.
    ///
    /// * For BUFFERED streams, data is made visible via a subsequent `FlushRows`
    ///   rpc which advances a cursor to a newer offset in the stream.
    ///
    /// * For PENDING streams, data is not made visible until the stream itself is
    ///   finalized (via the `FinalizeWriteStream` rpc), and the stream is explicitly
    ///   committed via the `BatchCommitWriteStreams` rpc.
    ///
    pub(crate) fn append_rows(&self) -> super::builder::big_query_write::AppendRows {
        super::builder::big_query_write::AppendRows::new(self.inner.clone())
    }

    /// Gets information about a write stream.
    pub(crate) fn get_write_stream(&self) -> super::builder::big_query_write::GetWriteStream {
        super::builder::big_query_write::GetWriteStream::new(self.inner.clone())
    }

    /// Finalize a write stream so that no new data can be appended to the
    /// stream. Finalize is not supported on the '_default' stream.
    pub(crate) fn finalize_write_stream(
        &self,
    ) -> super::builder::big_query_write::FinalizeWriteStream {
        super::builder::big_query_write::FinalizeWriteStream::new(self.inner.clone())
    }

    /// Atomically commits a group of `PENDING` streams that belong to the same
    /// `parent` table.
    ///
    /// Streams must be finalized before commit and cannot be committed multiple
    /// times. Once a stream is committed, data in the stream becomes available
    /// for read operations.
    pub(crate) fn batch_commit_write_streams(
        &self,
    ) -> super::builder::big_query_write::BatchCommitWriteStreams {
        super::builder::big_query_write::BatchCommitWriteStreams::new(self.inner.clone())
    }

    /// Flushes rows to a BUFFERED stream.
    ///
    /// If users are appending rows to BUFFERED stream, flush operation is
    /// required in order for the rows to become available for reading. A
    /// Flush operation flushes up to any previously flushed offset in a BUFFERED
    /// stream, to the offset specified in the request.
    ///
    /// Flush is not supported on the _default stream, since it is not BUFFERED.
    pub(crate) fn flush_rows(&self) -> super::builder::big_query_write::FlushRows {
        super::builder::big_query_write::FlushRows::new(self.inner.clone())
    }
}