use std::io::Read;
use std::sync::Arc;
use arrow_format::ipc;
use arrow_format::ipc::Schema::MetadataVersion;
use crate::array::*;
use crate::datatypes::Schema;
use crate::error::{ArrowError, Result};
use crate::record_batch::RecordBatch;
use super::super::convert;
use super::super::CONTINUATION_MARKER;
use super::common::*;
type ArrayRef = Arc<dyn Array>;
#[derive(Debug)]
pub struct StreamMetadata {
schema: Arc<Schema>,
version: MetadataVersion,
is_little_endian: bool,
}
pub fn read_stream_metadata<R: Read>(reader: &mut R) -> Result<StreamMetadata> {
let mut meta_size: [u8; 4] = [0; 4];
reader.read_exact(&mut meta_size)?;
let meta_len = {
if meta_size == CONTINUATION_MARKER {
reader.read_exact(&mut meta_size)?;
}
i32::from_le_bytes(meta_size)
};
let mut meta_buffer = vec![0; meta_len as usize];
reader.read_exact(&mut meta_buffer)?;
let message = ipc::Message::root_as_message(meta_buffer.as_slice())
.map_err(|err| ArrowError::Ipc(format!("Unable to get root as message: {:?}", err)))?;
let version = message.version();
let ipc_schema: ipc::Schema::Schema = message
.header_as_schema()
.ok_or_else(|| ArrowError::Ipc("Unable to read IPC message as schema".to_string()))?;
let (schema, is_little_endian) = convert::fb_to_schema(ipc_schema);
let schema = Arc::new(schema);
Ok(StreamMetadata {
schema,
version,
is_little_endian,
})
}
pub enum StreamState {
Waiting,
Some(RecordBatch),
}
impl StreamState {
pub fn unwrap(self) -> RecordBatch {
if let StreamState::Some(batch) = self {
batch
} else {
panic!("The batch is not available")
}
}
}
pub fn read_next<R: Read>(
reader: &mut R,
metadata: &StreamMetadata,
dictionaries_by_field: &mut Vec<Option<ArrayRef>>,
) -> Result<Option<StreamState>> {
let mut meta_size: [u8; 4] = [0; 4];
match reader.read_exact(&mut meta_size) {
Ok(()) => (),
Err(e) => {
return if e.kind() == std::io::ErrorKind::UnexpectedEof {
Ok(Some(StreamState::Waiting))
} else {
Err(ArrowError::from(e))
};
}
}
let meta_len = {
if meta_size == CONTINUATION_MARKER {
reader.read_exact(&mut meta_size)?;
}
i32::from_le_bytes(meta_size)
};
if meta_len == 0 {
return Ok(None);
}
let mut meta_buffer = vec![0; meta_len as usize];
reader.read_exact(&mut meta_buffer)?;
let vecs = &meta_buffer.to_vec();
let message = ipc::Message::root_as_message(vecs)
.map_err(|err| ArrowError::Ipc(format!("Unable to get root as message: {:?}", err)))?;
match message.header_type() {
ipc::Message::MessageHeader::Schema => Err(ArrowError::Ipc(
"Not expecting a schema when messages are read".to_string(),
)),
ipc::Message::MessageHeader::RecordBatch => {
let batch = message.header_as_record_batch().ok_or_else(|| {
ArrowError::Ipc("Unable to read IPC message as record batch".to_string())
})?;
let mut buf = vec![0; message.bodyLength() as usize];
reader.read_exact(&mut buf)?;
let mut reader = std::io::Cursor::new(buf);
read_record_batch(
batch,
metadata.schema.clone(),
None,
metadata.is_little_endian,
dictionaries_by_field,
metadata.version,
&mut reader,
0,
)
.map(|x| Some(StreamState::Some(x)))
}
ipc::Message::MessageHeader::DictionaryBatch => {
let batch = message.header_as_dictionary_batch().ok_or_else(|| {
ArrowError::Ipc("Unable to read IPC message as dictionary batch".to_string())
})?;
let mut buf = vec![0; message.bodyLength() as usize];
reader.read_exact(&mut buf)?;
let mut dict_reader = std::io::Cursor::new(buf);
read_dictionary(
batch,
&metadata.schema,
metadata.is_little_endian,
dictionaries_by_field,
&mut dict_reader,
0,
)?;
read_next(reader, metadata, dictionaries_by_field)
}
ipc::Message::MessageHeader::NONE => Ok(Some(StreamState::Waiting)),
t => Err(ArrowError::Ipc(format!(
"Reading types other than record batches not yet supported, unable to read {:?} ",
t
))),
}
}
pub struct StreamReader<R: Read> {
reader: R,
metadata: StreamMetadata,
dictionaries_by_field: Vec<Option<ArrayRef>>,
finished: bool,
}
impl<R: Read> StreamReader<R> {
pub fn new(reader: R, metadata: StreamMetadata) -> Self {
let fields = metadata.schema.fields().len();
Self {
reader,
metadata,
dictionaries_by_field: vec![None; fields],
finished: false,
}
}
pub fn schema(&self) -> &Arc<Schema> {
&self.metadata.schema
}
pub fn is_finished(&self) -> bool {
self.finished
}
fn maybe_next(&mut self) -> Result<Option<StreamState>> {
if self.finished {
return Ok(None);
}
let batch = read_next(
&mut self.reader,
&self.metadata,
&mut self.dictionaries_by_field,
)?;
if batch.is_none() {
self.finished = true;
}
Ok(batch)
}
}
impl<R: Read> Iterator for StreamReader<R> {
type Item = Result<StreamState>;
fn next(&mut self) -> Option<Self::Item> {
self.maybe_next().transpose()
}
}