use std::fs::File;
use std::path::{Path, PathBuf};
use polars::prelude::*;
use super::FrameSink;
use crate::RpoError;
pub struct ParquetSink {
path: PathBuf,
accumulated: Option<DataFrame>,
}
impl ParquetSink {
pub fn new(path: impl AsRef<Path>) -> Result<Self, RpoError> {
Ok(Self {
path: path.as_ref().to_path_buf(),
accumulated: None,
})
}
}
impl FrameSink for ParquetSink {
fn write_snapshot(&mut self, frame: DataFrame) -> Result<(), RpoError> {
match &mut self.accumulated {
None => {
self.accumulated = Some(frame);
}
Some(acc) => {
acc.vstack_mut(&frame)?;
}
}
Ok(())
}
fn finish(&mut self) -> Result<(), RpoError> {
let Some(mut df) = self.accumulated.take() else {
let mut empty = DataFrame::empty();
let file = File::create(&self.path)?;
ParquetWriter::new(file).finish(&mut empty)?;
return Ok(());
};
df.align_chunks_par();
let file = File::create(&self.path)?;
ParquetWriter::new(file).finish(&mut df)?;
Ok(())
}
}
pub struct ParquetDirSink {
dir: PathBuf,
}
impl ParquetDirSink {
pub fn new(dir: impl AsRef<Path>) -> Result<Self, RpoError> {
let p = dir.as_ref();
std::fs::create_dir_all(p)?;
Ok(Self {
dir: p.to_path_buf(),
})
}
}
impl FrameSink for ParquetDirSink {
fn write_snapshot(&mut self, mut frame: DataFrame) -> Result<(), RpoError> {
let sha = frame
.column("snapshot_sha")?
.str()?
.get(0)
.ok_or_else(|| RpoError::Sink("empty snapshot frame".into()))?
.to_string();
let path = self.dir.join(format!("{sha}.parquet"));
let file = File::create(&path)?;
ParquetWriter::new(file).finish(&mut frame)?;
Ok(())
}
fn finish(&mut self) -> Result<(), RpoError> {
Ok(())
}
}