Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions src/bigquery/src/write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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;
Expand Down
10 changes: 8 additions & 2 deletions src/bigquery/src/write/arrow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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:
/// <https://github.com/googleapis/google-cloud-rust/issues/6855>
pub type DefaultWriter = super::DefaultWriter<Arrow>;
8 changes: 6 additions & 2 deletions src/bigquery/src/write/arrow/writer_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
))
}

Expand Down Expand Up @@ -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(())
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Comment thread
dbolduc marked this conversation as resolved.
pub struct DefaultWriter {
pub struct DefaultWriter<F> {
pub(crate) inner: Arc<Dispatcher>,
pub(crate) write_stream: String,
pub(crate) schema: ArrowSchema,
pub(crate) format: F,
}

impl DefaultWriter {
impl<F> DefaultWriter<F>
where
F: DataFormat,
{
pub(crate) fn new(
pool: Arc<StreamPool>,
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;
Comment thread
dbolduc marked this conversation as resolved.
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<()> {
Comment thread
dbolduc marked this conversation as resolved.
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);
Expand All @@ -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?;
Expand Down
4 changes: 4 additions & 0 deletions src/bigquery/src/write/format.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand All @@ -31,3 +32,6 @@ pub(super) mod sealed {
Self: super::DataFormat;
}
}

pub use arrow::Arrow;
pub(crate) use proto::Proto;
6 changes: 4 additions & 2 deletions src/bigquery/src/write/proto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Proto>;
134 changes: 0 additions & 134 deletions src/bigquery/src/write/proto/default.rs

This file was deleted.

Loading
Loading