pub mod avro_types;
pub(crate) mod export_keys;
pub mod export_modal;
pub mod nested_json;
pub mod output_file;
pub mod python_script;
use std::io::{self, BufWriter, Write};
use std::path::PathBuf;
use std::time::{Duration, Instant};
use color_eyre::Result;
use polars::prelude::{
CsvWriter, DataFrame, IpcWriter, JsonFormat, JsonWriter, LazyFrame, ParquetWriter, SerWriter,
};
use crate::CompressionFormat;
use crate::export::export_modal::ExportFormat;
use crate::export::output_file::{OutputFile, Overwrite};
#[derive(Debug, Clone)]
pub struct ExportOptions {
pub csv_delimiter: u8,
pub csv_include_header: bool,
pub source_file: bool,
pub csv_compression: Option<CompressionFormat>,
pub json_compression: Option<CompressionFormat>,
pub ndjson_compression: Option<CompressionFormat>,
}
impl ExportOptions {
pub fn compression(&self, format: ExportFormat) -> Option<CompressionFormat> {
match format {
ExportFormat::Csv | ExportFormat::Tsv | ExportFormat::Psv => self.csv_compression,
ExportFormat::Json => self.json_compression,
ExportFormat::Ndjson => self.ndjson_compression,
ExportFormat::Parquet | ExportFormat::Ipc | ExportFormat::Avro => None,
}
}
pub fn delimiter(&self, format: ExportFormat) -> u8 {
format.preset_delimiter().unwrap_or(self.csv_delimiter)
}
}
#[derive(Debug, Clone)]
pub struct ExportRequest {
pub path: PathBuf,
pub format: ExportFormat,
pub options: ExportOptions,
pub overwrite: Overwrite,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Route {
Streamed,
Collected,
}
impl ExportRequest {
pub fn route(&self, polars_streaming: bool) -> Route {
let sinkable = (self.format.is_delimited() || self.format == ExportFormat::Parquet)
&& self.options.compression(self.format).is_none();
if cfg!(feature = "streaming") && polars_streaming && sinkable {
Route::Streamed
} else {
Route::Collected
}
}
}
const PROGRESS_EVERY: Duration = Duration::from_millis(250);
pub fn run(
lf: LazyFrame,
request: &ExportRequest,
polars_streaming: bool,
mut written: impl FnMut(u64) + Send + 'static,
) -> Result<()> {
let lf = request.format.prepare(lf)?;
let polars_streaming = crate::analysis::statistics::may_stream(&lf, polars_streaming);
let mut out = OutputFile::create(&request.path, request.overwrite)?;
match request.route(polars_streaming) {
#[cfg(feature = "streaming")]
Route::Streamed => {
written(0);
let file = Counted::new(out.file().try_clone()?, written);
sink(lf, request.format, &request.options, file.into_writable())?;
}
_ => {
let mut df = crate::analysis::statistics::collect_lazy(lf, polars_streaming)?;
written(0);
let file = Counted::new(out.file(), written);
encode(&mut df, request.format, &request.options, file)?;
}
}
out.commit()?;
Ok(())
}
#[cfg(feature = "streaming")]
fn sink(
lf: LazyFrame,
format: ExportFormat,
options: &ExportOptions,
writable: polars::io::utils::file::Writable,
) -> Result<()> {
use polars::prelude::{
CompatLevel, CsvWriterOptions, Engine, FileWriteFormat, ParquetWriteOptions,
SerializeOptions, SinkDestination, SinkTarget, SpecialEq, UnifiedSinkArgs,
};
use std::sync::{Arc, Mutex};
let file_format = match format {
ExportFormat::Csv | ExportFormat::Tsv | ExportFormat::Psv => {
FileWriteFormat::Csv(CsvWriterOptions {
include_header: options.csv_include_header,
serialize_options: Arc::new(SerializeOptions {
separator: options.delimiter(format),
..SerializeOptions::default()
}),
..CsvWriterOptions::default()
})
}
ExportFormat::Parquet => FileWriteFormat::Parquet(Arc::new(ParquetWriteOptions {
compat_level: Some(CompatLevel::newest()),
..ParquetWriteOptions::default()
})),
other => unreachable!("{other:?} does not stream"),
};
let target = SinkTarget::Dyn(SpecialEq::new(Arc::new(Mutex::new(Some(writable)))));
lf.sink(
SinkDestination::File { target },
file_format,
UnifiedSinkArgs::default(),
)?
.collect_with_engine(Engine::Streaming)?;
Ok(())
}
struct Counted<W, F> {
inner: W,
bytes: u64,
reported: Instant,
report: F,
}
impl<W: Write, F: FnMut(u64)> Counted<W, F> {
fn new(inner: W, report: F) -> Self {
Self {
inner,
bytes: 0,
reported: Instant::now(),
report,
}
}
}
impl<W: Write, F: FnMut(u64)> Write for Counted<W, F> {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let n = self.inner.write(buf)?;
self.bytes += n as u64;
if self.reported.elapsed() >= PROGRESS_EVERY {
self.reported = Instant::now();
(self.report)(self.bytes);
}
Ok(n)
}
fn flush(&mut self) -> io::Result<()> {
self.inner.flush()
}
}
#[cfg(feature = "streaming")]
impl<F: FnMut(u64) + Send + 'static> Counted<std::fs::File, F> {
fn into_writable(self) -> polars::io::utils::file::Writable {
polars::io::utils::file::Writable::Dyn(Box::new(self))
}
}
#[cfg(feature = "streaming")]
impl<F: FnMut(u64)> polars::io::utils::file::WritableTrait for Counted<std::fs::File, F> {
fn close(&mut self) -> io::Result<()> {
self.inner.flush()
}
fn sync_all(&self) -> io::Result<()> {
self.inner.sync_all()
}
fn sync_data(&self) -> io::Result<()> {
self.inner.sync_data()
}
}
pub fn encode<W: Write>(
df: &mut DataFrame,
format: ExportFormat,
options: &ExportOptions,
sink: W,
) -> Result<()> {
let mut sink = match options.compression(format) {
None => {
let mut buffered = BufWriter::new(sink);
serialize(df, format, options, &mut buffered)?;
buffered
.into_inner()
.map_err(io::IntoInnerError::into_error)?
}
Some(compression) => {
let mut buffered = BufWriter::new(Encoder::new(compression, sink)?);
serialize(df, format, options, &mut buffered)?;
let encoder = buffered
.into_inner()
.map_err(io::IntoInnerError::into_error)?;
encoder.finish()?
}
};
sink.flush()?;
Ok(())
}
fn serialize(
df: &mut DataFrame,
format: ExportFormat,
options: &ExportOptions,
out: &mut impl Write,
) -> Result<()> {
match format {
ExportFormat::Csv | ExportFormat::Tsv | ExportFormat::Psv => CsvWriter::new(out)
.with_separator(options.delimiter(format))
.include_header(options.csv_include_header)
.finish(df)?,
ExportFormat::Parquet => {
ParquetWriter::new(out).finish(df)?;
}
ExportFormat::Json => JsonWriter::new(out)
.with_json_format(JsonFormat::Json)
.finish(df)?,
ExportFormat::Ndjson => JsonWriter::new(out)
.with_json_format(JsonFormat::JsonLines)
.finish(df)?,
ExportFormat::Ipc => IpcWriter::new(out).finish(df)?,
ExportFormat::Avro => crate::export::avro_types::write(df, out)?,
}
Ok(())
}
enum Encoder<W: Write> {
Gzip(flate2::write::GzEncoder<W>),
Zstd(zstd::Encoder<'static, W>),
Bzip2(bzip2::write::BzEncoder<W>),
Xz(xz2::write::XzEncoder<W>),
}
impl<W: Write> Encoder<W> {
fn new(compression: CompressionFormat, out: W) -> io::Result<Self> {
Ok(match compression {
CompressionFormat::Gzip => Self::Gzip(flate2::write::GzEncoder::new(
out,
flate2::Compression::default(),
)),
CompressionFormat::Zstd => Self::Zstd(zstd::Encoder::new(out, 0)?),
CompressionFormat::Bzip2 => Self::Bzip2(bzip2::write::BzEncoder::new(
out,
bzip2::Compression::default(),
)),
CompressionFormat::Xz => Self::Xz(xz2::write::XzEncoder::new(out, 6)),
})
}
fn finish(self) -> io::Result<W> {
match self {
Self::Gzip(e) => e.finish(),
Self::Zstd(e) => e.finish(),
Self::Bzip2(e) => e.finish(),
Self::Xz(e) => e.finish(),
}
}
}
impl<W: Write> Write for Encoder<W> {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
match self {
Self::Gzip(e) => e.write(buf),
Self::Zstd(e) => e.write(buf),
Self::Bzip2(e) => e.write(buf),
Self::Xz(e) => e.write(buf),
}
}
fn flush(&mut self) -> io::Result<()> {
match self {
Self::Gzip(e) => e.flush(),
Self::Zstd(e) => e.flush(),
Self::Bzip2(e) => e.flush(),
Self::Xz(e) => e.flush(),
}
}
}
#[cfg(test)]
mod tests;