#[cfg(feature = "csv")]
mod csv;
#[cfg(feature = "parquet")]
mod parquet;
#[cfg(any(feature = "csv", feature = "parquet"))]
use crate::error::Error;
use crate::{
cell::CellValue,
dataset::DatasetMetadata,
error::Result,
parser::{ColumnInfo, ColumnarBatch, DatasetLayout, StreamingRow},
};
#[cfg(feature = "csv")]
pub use csv::CsvSink;
#[cfg(feature = "parquet")]
pub use parquet::ParquetSink;
#[cfg(any(feature = "csv", feature = "parquet"))]
use std::borrow::Cow;
pub struct SinkContext<'a> {
pub metadata: &'a DatasetMetadata,
pub columns: &'a [ColumnInfo],
pub source_path: Option<String>,
}
impl<'a> SinkContext<'a> {
#[must_use]
pub fn new(parsed: &'a DatasetLayout) -> Self {
Self {
metadata: &parsed.header.metadata,
columns: &parsed.columns,
source_path: None,
}
}
}
pub trait RowSink {
fn begin(&mut self, context: SinkContext<'_>) -> Result<()>;
fn write_row(&mut self, row: &[CellValue<'_>]) -> Result<()>;
fn write_streaming_row(&mut self, row: StreamingRow<'_, '_>) -> Result<()> {
let values = row.materialize()?;
self.write_row(&values)
}
fn finish(&mut self) -> Result<()>;
}
pub trait ColumnarSink: RowSink {
fn write_columnar_batch(
&mut self,
batch: &ColumnarBatch<'_>,
selection: &[usize],
) -> Result<()>;
}
#[cfg(any(feature = "csv", feature = "parquet"))]
pub(crate) fn validate_sink_begin(
context: &SinkContext<'_>,
writer_present: bool,
sink_name: &str,
) -> Result<()> {
if writer_present {
return Err(Error::Unsupported {
feature: Cow::Owned(format!(
"{sink_name} sink cannot be reused without finishing"
)),
});
}
if context.metadata.variables.len() != context.columns.len() {
return Err(Error::InvalidMetadata {
details: Cow::from("column metadata length mismatch"),
});
}
Ok(())
}