pub use wp_connector_utils::arrow::WireFormat;
use arrow::record_batch::RecordBatch;
use wf_connector_api::{SourceReason, SourceResult};
use wp_connector_api::SourceBatch;
use super::payload::payload_to_bytes;
pub fn decode_arrow_ipc_batches(events: &SourceBatch) -> SourceResult<Vec<RecordBatch>> {
let mut batches = Vec::new();
for event in events {
let payload = payload_to_bytes(&event.payload);
let decoded = wp_connector_utils::arrow::decode_arrow_ipc_batches(&payload)
.map_err(|e| SourceReason::Decode.err_detail(e))?;
batches.extend(decoded);
}
Ok(batches)
}
pub fn decode_arrow_framed_batches(events: &SourceBatch) -> SourceResult<Vec<RecordBatch>> {
let mut batches = Vec::new();
for event in events {
let payload = payload_to_bytes(&event.payload);
let decoded = wp_connector_utils::arrow::decode_arrow_framed_batches(&payload)
.map_err(|e| SourceReason::Decode.err_detail(e))?;
batches.extend(decoded);
}
Ok(batches)
}