Skip to main content

google_cloud_bigquery/write/
buffered.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::base::BaseWriter;
16use super::format::DataFormat;
17use crate::Result;
18use crate::model::{FinalizeWriteStreamResponse, FlushRowsResponse};
19use crate::write::builder::AppendWithOffset;
20use crate::write::transport::Transport;
21use std::sync::Arc;
22
23/// A writer for a [buffered stream].
24///
25/// In a buffered stream, row-level commits are provided, and records are
26/// buffered until the rows are committed by flushing the stream. This is an
27/// advanced stream type; if you have small batches that you want to guarantee
28/// appear together, consider using a [committed stream][crate::write::CommittedWriter]
29/// and sending each batch in one request.
30///
31/// [buffered stream]: https://docs.cloud.google.com/bigquery/docs/write-api-grpc#buffered_type
32#[derive(Debug)]
33pub struct BufferedWriter<F> {
34    pub(crate) inner: BaseWriter<F>,
35}
36
37impl<F> BufferedWriter<F>
38where
39    F: DataFormat,
40{
41    pub(crate) fn new(inner: Arc<Transport>, write_stream: String, format: F) -> Self {
42        Self {
43            inner: BaseWriter::new(inner, write_stream, format),
44        }
45    }
46
47    /// Returns the full resource name of the underlying write stream.
48    pub fn write_stream(&self) -> &str {
49        &self.inner.write_stream
50    }
51
52    /// Appends rows to the buffered stream.
53    pub fn append(&self, rows: F::Rows) -> AppendWithOffset {
54        AppendWithOffset::new(
55            self.inner.runner.req_tx.clone(),
56            self.inner.append_request(rows),
57        )
58    }
59
60    /// Flush the buffered stream, making rows up to and including the specified
61    /// offset available for reading.
62    ///
63    /// Stream offsets are 0-indexed and `offset` is **inclusive**. For example,
64    /// after appending a batch of 10 rows starting at offset 0 (occupying
65    /// offsets `0..=9`), calling `flush(9)` flushes all 10 rows.
66    ///
67    /// # Example
68    ///
69    /// ```
70    /// # use google_cloud_bigquery::write::BufferedWriter;
71    /// # use google_cloud_bigquery::write::format::Arrow;
72    /// # async fn sample(writer: BufferedWriter<Arrow>) -> anyhow::Result<()> {
73    /// // Append 10 rows starting at offset 0 (offsets 0..=9).
74    /// let _ = writer.append(ten_rows()).set_offset(0).send().await?;
75    ///
76    /// // Flush rows up to and including offset 9.
77    /// let _ = writer.flush(9).await?;
78    /// # Ok(()) }
79    ///
80    /// use google_cloud_bigquery::model::ArrowRecordBatch;
81    /// fn ten_rows() -> ArrowRecordBatch {
82    ///     todo!("Serialize 10 rows...")
83    /// }
84    /// ```
85    pub async fn flush(&self, offset: i64) -> Result<FlushRowsResponse> {
86        self.inner
87            .client
88            .flush_rows()
89            .set_write_stream(&self.inner.write_stream)
90            .set_offset(offset)
91            .send()
92            .await
93    }
94
95    /// Finalizes the buffered stream, preventing further writes.
96    pub async fn finalize(&self) -> Result<FinalizeWriteStreamResponse> {
97        self.inner.finalize().await
98    }
99}
100
101#[cfg(test)]
102mod tests {
103    use super::*;
104    use crate::error::AppendError;
105    use crate::write::test::*;
106    use bigquery_grpc_mock::{MockBigQueryWrite, start};
107    use gaxi::grpc::tonic::Response as TonicResponse;
108    use tokio::sync::mpsc;
109
110    #[tokio::test]
111    async fn basic_success() -> anyhow::Result<()> {
112        let (response_tx, response_rx) = mpsc::channel(10);
113
114        let mut mock = MockBigQueryWrite::new();
115        mock.expect_append_rows()
116            .return_once(|_| Ok(TonicResponse::from(response_rx)));
117
118        mock.expect_flush_rows()
119            .return_once(|req| {
120                assert_eq!(req.get_ref().offset, Some(3));
121                assert_eq!(req.get_ref().write_stream, write_stream());
122                Ok(TonicResponse::new(
123                    bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse::default()
124                ))
125            });
126
127        mock.expect_finalize_write_stream()
128            .return_once(|req| {
129                assert_eq!(req.get_ref().name, write_stream());
130                Ok(TonicResponse::new(
131                    bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FinalizeWriteStreamResponse::default()
132                ))
133            });
134
135        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
136        let transport = Arc::new(test_transport(endpoint).await?);
137
138        let writer = BufferedWriter::new(transport, write_stream(), format());
139        assert_eq!(writer.write_stream(), write_stream());
140
141        response_tx.send(Ok(convert(&test_response(1)))).await?;
142        let resp = writer.append(rows(1)).send().await?;
143        assert_eq!(resp.offset, Some(1));
144
145        response_tx.send(Ok(convert(&test_response(2)))).await?;
146        let resp = writer.append(rows(2)).send().await?;
147        assert_eq!(resp.offset, Some(2));
148
149        response_tx.send(Ok(convert(&test_response(3)))).await?;
150        let resp = writer.append(rows(3)).send().await?;
151        assert_eq!(resp.offset, Some(3));
152
153        drop(response_tx);
154        let err = writer.append(rows(4)).send().await.expect_err("channel");
155        assert!(matches!(err, AppendError::UnexpectedEndOfStream));
156
157        writer.flush(3).await?;
158        writer.finalize().await?;
159
160        Ok(())
161    }
162
163    #[tokio::test]
164    async fn multiple_flushes() -> anyhow::Result<()> {
165        let (response_tx, response_rx) = mpsc::channel(10);
166        let mut mock = MockBigQueryWrite::new();
167        mock.expect_append_rows()
168            .return_once(|_| Ok(TonicResponse::from(response_rx)));
169
170        mock.expect_flush_rows().times(2).returning(|req| {
171            Ok(TonicResponse::new(
172                bigquery_grpc_mock::google::cloud::bigquery::storage::v1::FlushRowsResponse {
173                    offset: req.get_ref().offset.unwrap_or(0),
174                },
175            ))
176        });
177
178        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
179        let transport = Arc::new(test_transport(endpoint).await?);
180        let writer = BufferedWriter::new(transport, write_stream(), format());
181        assert_eq!(writer.write_stream(), write_stream());
182
183        response_tx.send(Ok(convert(&test_response(1)))).await?;
184        let _ = writer.append(rows(1)).send().await?;
185        let flush1 = writer.flush(1).await?;
186        assert_eq!(flush1.offset, 1);
187
188        response_tx.send(Ok(convert(&test_response(2)))).await?;
189        let _ = writer.append(rows(2)).send().await?;
190        let flush2 = writer.flush(2).await?;
191        assert_eq!(flush2.offset, 2);
192
193        Ok(())
194    }
195}