pub mod batch;
pub mod schema;
pub use batch::{data_record_to_batch, data_records_to_batch};
pub use schema::{infer_arrow_schema, infer_schema_from_record};
use arrow::ipc::writer::StreamWriter;
use arrow::record_batch::RecordBatch;
use orion_error::conversion::{SourceRawErr, ToStructError};
use wp_connector_api::{SinkReason, SinkResult};
pub(crate) fn sink_err<E>(msg: &'static str, err: E) -> wp_connector_api::SinkError
where
E: std::fmt::Display,
{
SinkReason::Sink
.to_err()
.with_detail(format!("{msg}: {err}"))
}
pub(crate) fn encode_batch_ipc_stream(batch: &RecordBatch) -> SinkResult<Vec<u8>> {
let schema = batch.schema();
let mut buf = Vec::new();
{
let mut writer = StreamWriter::try_new(&mut buf, &schema)
.source_raw_err(SinkReason::Sink, "arrow create stream writer")?;
writer
.write(batch)
.map_err(|e| sink_err("arrow encode batch", e))?;
writer
.finish()
.map_err(|e| sink_err("arrow finish stream", e))?;
}
Ok(buf)
}
pub(crate) fn encode_ipc_frame(tag: &str, batch: &RecordBatch) -> SinkResult<Vec<u8>> {
let tag_bytes = tag.as_bytes();
let schema = batch.schema();
let mut buf = Vec::with_capacity(4 + tag_bytes.len() + 1024);
buf.extend_from_slice(&(tag_bytes.len() as u32).to_be_bytes());
buf.extend_from_slice(tag_bytes);
{
let mut writer = StreamWriter::try_new(&mut buf, &schema)
.source_raw_err(SinkReason::Sink, "arrow create framed stream writer")?;
writer
.write(batch)
.map_err(|e| sink_err("arrow encode framed batch", e))?;
writer
.finish()
.map_err(|e| sink_err("arrow finish framed stream", e))?;
}
Ok(buf)
}
pub(crate) fn encode_ipc_frame_multi(tag: &str, batches: &[RecordBatch]) -> SinkResult<Vec<u8>> {
if batches.is_empty() {
let tag_bytes = tag.as_bytes();
let mut buf = Vec::with_capacity(4 + tag_bytes.len());
buf.extend_from_slice(&(tag_bytes.len() as u32).to_be_bytes());
buf.extend_from_slice(tag_bytes);
return Ok(buf);
}
let tag_bytes = tag.as_bytes();
let schema = batches[0].schema();
let mut buf = Vec::with_capacity(4 + tag_bytes.len() + 1024);
buf.extend_from_slice(&(tag_bytes.len() as u32).to_be_bytes());
buf.extend_from_slice(tag_bytes);
{
let mut writer = StreamWriter::try_new(&mut buf, &schema)
.source_raw_err(SinkReason::Sink, "arrow create framed multi stream writer")?;
for batch in batches {
writer
.write(batch)
.map_err(|e| sink_err("arrow encode framed multi batch", e))?;
}
writer
.finish()
.map_err(|e| sink_err("arrow finish framed multi stream", e))?;
}
Ok(buf)
}