use std::fs::File;
use std::io::BufWriter;
use std::path::{Path, PathBuf};
use arrow::ipc::reader::{FileReader, StreamReader};
use arrow::ipc::writer::{FileWriter, StreamWriter};
use arrow::record_batch::RecordBatch;
use log::{debug, info};
use crate::{StorageError, StorageResult};
pub const FILE_EXT: &str = "arrow";
pub const STREAM_EXT: &str = "arrows";
pub const FILE_STORAGE_FORMAT: &str = "arrow-ipc";
pub const STREAM_STORAGE_FORMAT: &str = "arrow-ipc-stream";
fn ipc_err(what: &str, e: arrow::error::ArrowError) -> StorageError {
StorageError::IPC(format!("{what}: {e}"))
}
pub fn write_ipc_file(path: &Path, batches: &[RecordBatch]) -> StorageResult<()> {
let batch = batches
.first()
.ok_or_else(|| StorageError::Invalid("cannot write an empty IPC file".into()))?;
let file = File::create(path).map_err(|e| {
StorageError::Io(format!("Failed to create IPC file {}: {e}", path.display()))
})?;
let mut writer = FileWriter::try_new_buffered(file, batch.schema().as_ref()).map_err(|e| {
ipc_err(
&format!("Failed to create IPC file writer at {}", path.display()),
e,
)
})?;
for batch in batches {
writer.write(batch).map_err(|e| {
ipc_err(
&format!("Failed to write IPC batch to {}", path.display()),
e,
)
})?;
}
writer
.finish()
.map_err(|e| ipc_err(&format!("Failed to finish IPC file {}", path.display()), e))?;
info!(
"Wrote Arrow IPC file {} ({} batches)",
path.display(),
batches.len()
);
Ok(())
}
pub fn read_ipc_file(path: &Path) -> StorageResult<RecordBatch> {
let file = File::open(path).map_err(|e| {
StorageError::Io(format!("Failed to open IPC file {}: {e}", path.display()))
})?;
let reader = FileReader::try_new(file, None).map_err(|e| {
ipc_err(
&format!("Failed to open IPC file reader for {}", path.display()),
e,
)
})?;
let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().map_err(|e| {
ipc_err(
&format!("Failed to read IPC file {}: {e}", path.display()),
e,
)
})?;
combined_batches(path, batches)
}
pub fn read_ipc_stream(path: &Path) -> StorageResult<RecordBatch> {
let file = File::open(path).map_err(|e| {
StorageError::Io(format!("Failed to open IPC stream {}: {e}", path.display()))
})?;
let reader = StreamReader::try_new(file, None).map_err(|e| {
ipc_err(
&format!(
"Failed to open IPC stream reader for {}: {e}",
path.display()
),
e,
)
})?;
let batches: Vec<RecordBatch> = reader.collect::<Result<Vec<_>, _>>().map_err(|e| {
ipc_err(
&format!("Failed to read IPC stream {}: {e}", path.display()),
e,
)
})?;
combined_batches(path, batches)
}
pub fn read_ipc_artifact(path: &Path) -> StorageResult<RecordBatch> {
match path.extension().and_then(|e| e.to_str()) {
Some(FILE_EXT) => read_ipc_file(path),
Some(STREAM_EXT) => read_ipc_stream(path),
other => Err(StorageError::Invalid(format!(
"not an Arrow IPC artifact extension '{other:?}': {}",
path.display()
))),
}
}
fn combined_batches(path: &Path, batches: Vec<RecordBatch>) -> StorageResult<RecordBatch> {
if batches.is_empty() {
return Err(StorageError::Invalid(format!(
"Empty Arrow IPC artifact at {}",
path.display()
)));
}
let schema = batches[0].schema();
let combined = arrow::compute::concat_batches(&schema, &batches).map_err(|e| {
ipc_err(
&format!("Failed to concatenate IPC batches at {}", path.display()),
e,
)
})?;
debug!(
"IPC artifact {}: {} batches -> {} rows",
path.display(),
batches.len(),
combined.num_rows()
);
Ok(combined)
}
pub struct IpcStreamWriterHandle {
writer: Option<StreamWriter<BufWriter<File>>>,
path: PathBuf,
}
impl IpcStreamWriterHandle {
pub fn create(path: PathBuf, schema: &arrow::datatypes::Schema) -> StorageResult<Self> {
let file = File::create(&path).map_err(|e| {
StorageError::Io(format!(
"Failed to create IPC stream {}: {e}",
path.display()
))
})?;
let writer = StreamWriter::try_new_buffered(file, schema).map_err(|e| {
ipc_err(
&format!("Failed to create IPC stream writer at {}", path.display()),
e,
)
})?;
info!("Opened Arrow IPC stream writer at {}", path.display());
Ok(Self {
writer: Some(writer),
path,
})
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn write_batch(&mut self, batch: &RecordBatch) -> StorageResult<()> {
let writer = self
.writer
.as_mut()
.ok_or_else(|| StorageError::Invalid("IPC stream writer is already finished".into()))?;
writer.write(batch).map_err(|e| {
ipc_err(
&format!(
"Failed to write IPC stream batch to {}",
self.path.display()
),
e,
)
})
}
pub fn finish(mut self) -> StorageResult<PathBuf> {
let mut writer = self
.writer
.take()
.ok_or_else(|| StorageError::Invalid("IPC stream writer is already finished".into()))?;
writer.finish().map_err(|e| {
ipc_err(
&format!("Failed to finish IPC stream at {}", self.path.display()),
e,
)
})?;
info!("Closed Arrow IPC stream at {}", self.path.display());
Ok(self.path)
}
}
pub async fn write_ipc_file_async(path: PathBuf, batches: Vec<RecordBatch>) -> StorageResult<()> {
tokio::task::spawn_blocking(move || write_ipc_file(&path, &batches))
.await
.map_err(|e| StorageError::Io(format!("ipc writer task failed: {e}")))?
}
pub async fn read_ipc_artifact_async(path: PathBuf) -> StorageResult<RecordBatch> {
tokio::task::spawn_blocking(move || read_ipc_artifact(&path))
.await
.map_err(|e| StorageError::Io(format!("ipc reader task failed: {e}")))?
}
pub async fn open_stream_writer_async(
path: PathBuf,
schema: arrow::datatypes::SchemaRef,
) -> StorageResult<IpcStreamWriterHandle> {
tokio::task::spawn_blocking(move || IpcStreamWriterHandle::create(path, schema.as_ref()))
.await
.map_err(|e| StorageError::Io(format!("ipc stream writer task failed: {e}")))?
}