From 466361756e0eaf384f9eb75ba689b4e58ff58a64 Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 18 Sep 2026 15:48:27 -0400 Subject: [PATCH 1/4] refactor(bigquery): genericize `DefaultWriter` --- src/bigquery/src/write.rs | 8 ++- src/bigquery/src/write/arrow.rs | 5 +- .../src/write/arrow/writer_builder.rs | 8 ++- src/bigquery/src/write/{arrow => }/default.rs | 67 ++++++------------- src/bigquery/src/write/format.rs | 3 + src/bigquery/src/write/proto.rs | 4 +- .../src/write/proto/writer_builder.rs | 8 ++- 7 files changed, 47 insertions(+), 56 deletions(-) rename src/bigquery/src/write/{arrow => }/default.rs (59%) diff --git a/src/bigquery/src/write.rs b/src/bigquery/src/write.rs index d6228e8f24..3c5a216009 100644 --- a/src/bigquery/src/write.rs +++ b/src/bigquery/src/write.rs @@ -21,6 +21,11 @@ pub(crate) mod proto; pub use append_future::AppendFuture; +pub use default::DefaultWriter; + +/// Defines the data formats accepted by a writer. +pub mod format; + /// Defines the retry policy for the BigQuery Storage Write API. pub mod retry_policy; @@ -31,10 +36,9 @@ pub(super) mod client; pub(super) mod client_builder; pub(super) mod error; +mod default; mod dispatcher; mod entry; -#[cfg_attr(not(test), expect(dead_code))] -mod format; mod pool; mod proto_schema; mod runner; diff --git a/src/bigquery/src/write/arrow.rs b/src/bigquery/src/write/arrow.rs index 6db80c2ffd..42b4223be4 100644 --- a/src/bigquery/src/write/arrow.rs +++ b/src/bigquery/src/write/arrow.rs @@ -15,14 +15,15 @@ mod base; mod buffered; mod committed; -mod default; mod pending; mod writer; mod writer_builder; +use super::format::Arrow; + pub use buffered::BufferedWriter; pub use committed::CommittedWriter; -pub use default::DefaultWriter; +pub type DefaultWriter = super::DefaultWriter; pub use pending::PendingWriter; pub use writer::Writer; pub use writer_builder::WriterBuilder; diff --git a/src/bigquery/src/write/arrow/writer_builder.rs b/src/bigquery/src/write/arrow/writer_builder.rs index b45744c448..047f9e0ec6 100644 --- a/src/bigquery/src/write/arrow/writer_builder.rs +++ b/src/bigquery/src/write/arrow/writer_builder.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use super::super::format::Arrow; use super::super::generated::gapic_storage::client::BigQueryWrite; use super::super::pool::{StreamPool, StreamPoolOptions}; use super::super::retry_policy::RetryOptions; @@ -86,11 +87,14 @@ impl WriterBuilder { }; Arc::new(StreamPool::new(self.inner, options)) }; + let format = Arrow { + schema: self.schema, + }; Ok(DefaultWriter::new( pool, self.retry_options, write_stream, - self.schema, + format, )) } @@ -384,7 +388,7 @@ mod tests { writer.write_stream, "projects/p/datasets/d/tables/t/streams/_default" ); - assert_eq!(writer.schema, schema()); + assert_eq!(writer.format.schema, schema()); Ok(()) } diff --git a/src/bigquery/src/write/arrow/default.rs b/src/bigquery/src/write/default.rs similarity index 59% rename from src/bigquery/src/write/arrow/default.rs rename to src/bigquery/src/write/default.rs index 1a4881999e..1c1226a36a 100644 --- a/src/bigquery/src/write/arrow/default.rs +++ b/src/bigquery/src/write/default.rs @@ -12,88 +12,60 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::super::builder::Append; -use super::super::dispatcher::Dispatcher; -use super::super::pool::StreamPool; -use super::super::retry_policy::RetryOptions; -use crate::model::append_rows_request::ArrowData; -use crate::model::{AppendRowsRequest, ArrowRecordBatch, ArrowSchema}; +use super::builder::Append; +use super::dispatcher::Dispatcher; +use super::format::DataFormat; +use super::pool::StreamPool; +use super::retry_policy::RetryOptions; use std::sync::Arc; /// A writer for the [default stream]. /// /// [default stream]: https://docs.cloud.google.com/bigquery/docs/write-api#default_stream #[derive(Debug)] -pub struct DefaultWriter { +pub struct DefaultWriter { pub(crate) inner: Arc, pub(crate) write_stream: String, - pub(crate) schema: ArrowSchema, + pub(crate) format: F, } -impl DefaultWriter { +impl DefaultWriter +where + F: DataFormat, +{ pub(crate) fn new( pool: Arc, retry_options: RetryOptions, write_stream: String, - schema: ArrowSchema, + format: F, ) -> Self { let inner = Arc::new(Dispatcher::new(pool, retry_options)); Self { inner, write_stream, - schema, + format, } } /// Append rows to the stream. - pub fn append(&self, rows: ArrowRecordBatch) -> Append { - // TODO(#5744) - send optimization - let req = AppendRowsRequest::new() - .set_write_stream(&self.write_stream) - .set_arrow_rows( - ArrowData::new() - .set_writer_schema(self.schema.clone()) - .set_rows(rows), - ); + pub fn append(&self, rows: F::Rows) -> Append { + let req = self.format.make_request(&self.write_stream, rows); Append::new(self.inner.clone(), req) } } #[cfg(test)] mod tests { - use super::super::super::pool::StreamPoolOptions; + use super::super::format::Arrow; + use super::super::pool::StreamPoolOptions; use super::*; use crate::error::AppendError; + use crate::model::ArrowRecordBatch; use crate::write::test::*; use bigquery_grpc_mock::{MockBigQueryWrite, start}; use gaxi::grpc::tonic::{Response as TonicResponse, Status as TonicStatus}; use tokio::sync::mpsc; - #[tokio::test] - async fn request_fields() -> anyhow::Result<()> { - let transport = Arc::new(test_transport("http://ignored:1").await?); - let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default())); - let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), schema()); - - let b = writer.append(rows(1)); - assert_eq!(b.req.write_stream, write_stream()); - let data = b.req.arrow_rows().expect("arrow rows should be set"); - let s = data.writer_schema.as_ref().expect("schema should be set"); - assert_eq!(s.serialized_schema, "test"); - let r = data.rows.as_ref().expect("rows should be set"); - assert_eq!(r.serialized_record_batch, "1"); - - let b = writer.append(rows(2)); - assert_eq!(b.req.write_stream, write_stream()); - let data = b.req.arrow_rows().expect("arrow rows should be set"); - let s = data.writer_schema.as_ref().expect("schema should be set"); - assert_eq!(s.serialized_schema, "test"); - let r = data.rows.as_ref().expect("rows should be set"); - assert_eq!(r.serialized_record_batch, "2"); - - Ok(()) - } - #[tokio::test] async fn basic_success() -> anyhow::Result<()> { let (response_tx, response_rx) = mpsc::channel(10); @@ -105,7 +77,8 @@ mod tests { let transport = Arc::new(test_transport(endpoint).await?); let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default())); - let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), schema()); + let format = Arrow { schema: schema() }; + let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), format); response_tx.send(Ok(convert(&test_response(1)))).await?; let resp = writer.append(rows(1)).send().await?; diff --git a/src/bigquery/src/write/format.rs b/src/bigquery/src/write/format.rs index f67a661c56..472f4e3806 100644 --- a/src/bigquery/src/write/format.rs +++ b/src/bigquery/src/write/format.rs @@ -31,3 +31,6 @@ pub(super) mod sealed { Self: super::DataFormat; } } + +pub use arrow::Arrow; +pub(crate) use proto::Proto; diff --git a/src/bigquery/src/write/proto.rs b/src/bigquery/src/write/proto.rs index 95a23c18ba..66235357fc 100644 --- a/src/bigquery/src/write/proto.rs +++ b/src/bigquery/src/write/proto.rs @@ -20,9 +20,11 @@ mod pending; mod writer; mod writer_builder; +use super::format::Proto; + pub(crate) use buffered::BufferedWriter; pub(crate) use committed::CommittedWriter; -pub(crate) use default::DefaultWriter; +pub type DefaultWriter = super::DefaultWriter; pub(crate) use pending::PendingWriter; pub(crate) use writer::Writer; pub(crate) use writer_builder::WriterBuilder; diff --git a/src/bigquery/src/write/proto/writer_builder.rs b/src/bigquery/src/write/proto/writer_builder.rs index 6a9dab7b74..97a05b7e6f 100644 --- a/src/bigquery/src/write/proto/writer_builder.rs +++ b/src/bigquery/src/write/proto/writer_builder.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use super::super::format::Proto; use super::super::generated::gapic_storage::client::BigQueryWrite; use super::super::pool::{StreamPool, StreamPoolOptions}; use super::super::retry_policy::RetryOptions; @@ -61,11 +62,14 @@ impl WriterBuilder { ..Default::default() }; let pool = Arc::new(StreamPool::new(self.inner, options)); + let format = Proto { + schema: self.schema, + }; Ok(DefaultWriter::new( pool, self.retry_options, write_stream, - self.schema, + format, )) } @@ -268,7 +272,7 @@ mod tests { writer.write_stream, "projects/p/datasets/d/tables/t/streams/_default" ); - assert_eq!(writer.schema, proto_schema()); + assert_eq!(writer.format.schema, proto_schema()); Ok(()) } From 87074510e250c91c4eeb42038c43ca59358f5f75 Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 18 Sep 2026 16:58:58 -0400 Subject: [PATCH 2/4] fix public docs --- src/bigquery/src/write/arrow.rs | 7 ++++++- src/bigquery/src/write/format.rs | 1 + 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/src/bigquery/src/write/arrow.rs b/src/bigquery/src/write/arrow.rs index 42b4223be4..840a50e9f3 100644 --- a/src/bigquery/src/write/arrow.rs +++ b/src/bigquery/src/write/arrow.rs @@ -23,7 +23,12 @@ use super::format::Arrow; pub use buffered::BufferedWriter; pub use committed::CommittedWriter; -pub type DefaultWriter = super::DefaultWriter; pub use pending::PendingWriter; pub use writer::Writer; pub use writer_builder::WriterBuilder; + +/// DEPRECATED - do not use. +/// +/// This type is about to be deleted. See: +/// https://github.com/googleapis/google-cloud-rust/issues/6855 +pub type DefaultWriter = super::DefaultWriter; diff --git a/src/bigquery/src/write/format.rs b/src/bigquery/src/write/format.rs index 472f4e3806..5c514343e4 100644 --- a/src/bigquery/src/write/format.rs +++ b/src/bigquery/src/write/format.rs @@ -19,6 +19,7 @@ mod proto; /// /// This trait is sealed and cannot be implemented for types outside this crate. pub trait DataFormat: sealed::DataFormat { + /// The representation of rows for this data format. type Rows; } From 12d43b35ccf29c376d6b6c52ca0245431b1b8d1a Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 18 Sep 2026 17:10:51 -0400 Subject: [PATCH 3/4] fix docs --- src/bigquery/src/write/arrow.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/bigquery/src/write/arrow.rs b/src/bigquery/src/write/arrow.rs index 840a50e9f3..351d6c3fdd 100644 --- a/src/bigquery/src/write/arrow.rs +++ b/src/bigquery/src/write/arrow.rs @@ -30,5 +30,5 @@ pub use writer_builder::WriterBuilder; /// DEPRECATED - do not use. /// /// This type is about to be deleted. See: -/// https://github.com/googleapis/google-cloud-rust/issues/6855 +/// pub type DefaultWriter = super::DefaultWriter; From d86d1a5e2d10a749d7e351a3928598b46471cab9 Mon Sep 17 00:00:00 2001 From: Darren Bolduc Date: Fri, 18 Sep 2026 17:33:37 -0400 Subject: [PATCH 4/4] consistency; plus delete unused file --- src/bigquery/src/write.rs | 2 + src/bigquery/src/write/proto.rs | 4 +- src/bigquery/src/write/proto/default.rs | 134 ------------------------ 3 files changed, 4 insertions(+), 136 deletions(-) delete mode 100644 src/bigquery/src/write/proto/default.rs diff --git a/src/bigquery/src/write.rs b/src/bigquery/src/write.rs index 3c5a216009..3263a45a3e 100644 --- a/src/bigquery/src/write.rs +++ b/src/bigquery/src/write.rs @@ -12,10 +12,12 @@ // See the License for the specific language governing permissions and // limitations under the License. +// TODO(#6855) - delete this module /// Types to write data in [Arrow] format. /// /// [arrow]: https://arrow.apache.org/ pub mod arrow; +// TODO(#6855) - delete this module #[allow(dead_code)] pub(crate) mod proto; diff --git a/src/bigquery/src/write/proto.rs b/src/bigquery/src/write/proto.rs index 66235357fc..16e4b88217 100644 --- a/src/bigquery/src/write/proto.rs +++ b/src/bigquery/src/write/proto.rs @@ -15,7 +15,6 @@ mod base; mod buffered; mod committed; -mod default; mod pending; mod writer; mod writer_builder; @@ -24,7 +23,8 @@ use super::format::Proto; pub(crate) use buffered::BufferedWriter; pub(crate) use committed::CommittedWriter; -pub type DefaultWriter = super::DefaultWriter; pub(crate) use pending::PendingWriter; pub(crate) use writer::Writer; pub(crate) use writer_builder::WriterBuilder; + +pub(crate) type DefaultWriter = super::DefaultWriter; diff --git a/src/bigquery/src/write/proto/default.rs b/src/bigquery/src/write/proto/default.rs deleted file mode 100644 index c68ea806bd..0000000000 --- a/src/bigquery/src/write/proto/default.rs +++ /dev/null @@ -1,134 +0,0 @@ -// 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. - -use super::super::builder::Append; -use super::super::dispatcher::Dispatcher; -use super::super::pool::StreamPool; -use super::super::retry_policy::RetryOptions; -use crate::model::append_rows_request::ProtoData; -use crate::model::{AppendRowsRequest, ProtoRows, ProtoSchema}; -use std::sync::Arc; - -/// A writer for the [default stream] using Protobuf as the data format. -/// -/// [default stream]: https://docs.cloud.google.com/bigquery/docs/write-api#default_stream -#[derive(Debug)] -pub struct DefaultWriter { - inner: Arc, - pub(crate) write_stream: String, - pub(crate) schema: ProtoSchema, -} - -impl DefaultWriter { - pub(crate) fn new( - pool: Arc, - retry_options: RetryOptions, - write_stream: String, - schema: ProtoSchema, - ) -> Self { - let inner = Arc::new(Dispatcher::new(pool, retry_options)); - Self { - inner, - write_stream, - schema, - } - } - - /// Append rows to the stream. - pub fn append(&self, rows: ProtoRows) -> Append { - // TODO(#5744) - send optimization - let req = AppendRowsRequest::new() - .set_write_stream(&self.write_stream) - .set_proto_rows( - ProtoData::new() - .set_writer_schema(self.schema.clone()) - .set_rows(rows), - ); - Append::new(self.inner.clone(), req) - } -} - -#[cfg(test)] -mod tests { - use super::super::super::pool::StreamPoolOptions; - use super::*; - use crate::error::AppendError; - use crate::write::test::*; - use bigquery_grpc_mock::{MockBigQueryWrite, start}; - use gaxi::grpc::tonic::{Response as TonicResponse, Status as TonicStatus}; - use tokio::sync::mpsc; - - #[tokio::test] - async fn request_fields() -> anyhow::Result<()> { - let transport = Arc::new(test_transport("http://ignored:1").await?); - let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default())); - let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), proto_schema()); - - let b = writer.append(rows(1)); - assert_eq!(b.req.write_stream, write_stream()); - let data = b.req.proto_rows().expect("proto rows should be set"); - let s = data.writer_schema.as_ref().expect("schema should be set"); - assert_eq!(s.proto_descriptor.as_ref().unwrap().name, "TestMessage"); - let r = data.rows.as_ref().expect("rows should be set"); - assert_eq!(r.serialized_rows, vec![bytes::Bytes::from("1")]); - - let b = writer.append(rows(2)); - assert_eq!(b.req.write_stream, write_stream()); - let data = b.req.proto_rows().expect("proto rows should be set"); - let s = data.writer_schema.as_ref().expect("schema should be set"); - assert_eq!(s.proto_descriptor.as_ref().unwrap().name, "TestMessage"); - let r = data.rows.as_ref().expect("rows should be set"); - assert_eq!(r.serialized_rows, vec![bytes::Bytes::from("2")]); - - Ok(()) - } - - #[tokio::test] - async fn basic_success() -> anyhow::Result<()> { - let (response_tx, response_rx) = mpsc::channel(10); - - let mut mock = MockBigQueryWrite::new(); - mock.expect_append_rows() - .return_once(|_| Ok(TonicResponse::from(response_rx))); - let (endpoint, _server) = start("0.0.0.0:0", mock).await?; - let transport = Arc::new(test_transport(endpoint).await?); - let pool = Arc::new(StreamPool::new(transport, StreamPoolOptions::default())); - - let writer = DefaultWriter::new(pool, test_retry_options(), write_stream(), proto_schema()); - - response_tx.send(Ok(convert(&test_response(1)))).await?; - let resp = writer.append(rows(1)).send().await?; - assert_eq!(resp.offset, Some(1)); - - response_tx.send(Ok(convert(&test_response(2)))).await?; - let resp = writer.append(rows(2)).send().await?; - assert_eq!(resp.offset, Some(2)); - - response_tx.send(Ok(convert(&test_response(3)))).await?; - let resp = writer.append(rows(3)).send().await?; - assert_eq!(resp.offset, Some(3)); - - response_tx - .send(Err(TonicStatus::failed_precondition("fail"))) - .await?; - let err = writer.append(rows(4)).send().await.expect_err("fail"); - assert!(matches!(err, AppendError::Rpc { source: _ }), "{err:?}"); - - Ok(()) - } - - fn rows(id: i64) -> ProtoRows { - ProtoRows::new().set_serialized_rows(vec![bytes::Bytes::from(id.to_string())]) - } -}