Skip to content

Commit

Permalink
feat: Make AsyncArrowWriter accepts AsyncFileWriter
Browse files Browse the repository at this point in the history
Signed-off-by: Xuanwo <[email protected]>
  • Loading branch information
Xuanwo committed May 11, 2024
1 parent 68ecc16 commit 4d0fe42
Showing 1 changed file with 51 additions and 10 deletions.
61 changes: 51 additions & 10 deletions parquet/src/arrow/async_writer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,13 +59,59 @@ use crate::{
};
use arrow_array::RecordBatch;
use arrow_schema::SchemaRef;
use bytes::Bytes;
use futures::future::BoxFuture;
use futures::FutureExt;
use tokio::io::{AsyncWrite, AsyncWriteExt};

/// Encodes [`RecordBatch`] to parquet, outputting to an [`AsyncWrite`]
/// The asynchronous interface used by [`AsyncArrowWriter`] to write parquet files
pub trait AsyncFileWriter: Send {
/// Write the provided bytes to the underlying writer
///
/// The underlying writer CAN decide to buffer the data or write it immediately.
/// This design allows the writer implementer to control the buffering and I/O scheduling.
fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>>;

/// Flush any buffered data to the underlying writer and finish writing process.
///
/// After `complete` returns `Ok(())`, caller SHOULD not call write again.
fn complete(&mut self) -> BoxFuture<'_, Result<()>>;
}

impl AsyncFileWriter for Box<dyn AsyncFileWriter> {
fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>> {
self.as_mut().write(bs)
}

fn complete(&mut self) -> BoxFuture<'_, Result<()>> {
self.as_mut().complete()
}
}

impl<T: AsyncWrite + Unpin + Send> AsyncFileWriter for T {
fn write(&mut self, bs: Bytes) -> BoxFuture<'_, Result<()>> {
async move {
self.write_all(&bs).await?;
Ok(())
}
.boxed()
}

fn complete(&mut self) -> BoxFuture<'_, Result<()>> {
async move {
self.flush().await?;
self.shutdown().await?;
Ok(())
}
.boxed()
}
}

/// Encodes [`RecordBatch`] to parquet, outputting to an [`AsyncFileWriter`]
///
/// ## Memory Usage
///
/// This writer eagerly writes data as soon as possible to the underlying [`AsyncWrite`],
/// This writer eagerly writes data as soon as possible to the underlying [`AsyncFileWriter`],
/// permitting fine-grained control over buffering and I/O scheduling. However, the columnar
/// nature of parquet forces data for an entire row group to be buffered in memory, before
/// it can be flushed. Depending on the data and the configured row group size, this buffering
Expand Down Expand Up @@ -97,7 +143,7 @@ pub struct AsyncArrowWriter<W> {
async_writer: W,
}

impl<W: AsyncWrite + Unpin + Send> AsyncArrowWriter<W> {
impl<W: AsyncFileWriter> AsyncArrowWriter<W> {
/// Try to create a new Async Arrow Writer
pub fn try_new(
writer: W,
Expand Down Expand Up @@ -178,7 +224,7 @@ impl<W: AsyncWrite + Unpin + Send> AsyncArrowWriter<W> {

// Force to flush the remaining data.
self.do_write().await?;
self.async_writer.shutdown().await?;
self.async_writer.complete().await?;

Ok(metadata)
}
Expand All @@ -188,12 +234,7 @@ impl<W: AsyncWrite + Unpin + Send> AsyncArrowWriter<W> {
let buffer = self.sync_writer.inner_mut();

self.async_writer
.write_all(buffer.as_slice())
.await
.map_err(|e| ParquetError::External(Box::new(e)))?;

self.async_writer
.flush()
.write(Bytes::copy_from_slice(buffer))
.await
.map_err(|e| ParquetError::External(Box::new(e)))?;

Expand Down

0 comments on commit 4d0fe42

Please sign in to comment.