diff --git a/src/bigquery/src/write.rs b/src/bigquery/src/write.rs index d6228e8f24..3263a45a3e 100644 --- a/src/bigquery/src/write.rs +++ b/src/bigquery/src/write.rs @@ -12,15 +12,22 @@ // 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; 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 +38,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..351d6c3fdd 100644 --- a/src/bigquery/src/write/arrow.rs +++ b/src/bigquery/src/write/arrow.rs @@ -15,14 +15,20 @@ 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 use pending::PendingWriter; pub use writer::Writer; pub use writer_builder::WriterBuilder; + +/// DEPRECATED - do not use. +/// +/// This type is about to be deleted. See: +/// +pub type DefaultWriter = super::DefaultWriter; 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..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; } @@ -31,3 +32,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..16e4b88217 100644 --- a/src/bigquery/src/write/proto.rs +++ b/src/bigquery/src/write/proto.rs @@ -15,14 +15,16 @@ mod base; mod buffered; mod committed; -mod default; 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(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())]) - } -} 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(()) }