use lexical_core::ToLexical;
use crate::temporal_conversions;
use crate::types::{Index, NativeType};
use crate::util::lexical_to_bytes_mut;
use crate::{
array::{Array, BinaryArray, BooleanArray, PrimitiveArray, Utf8Array},
datatypes::{DataType, TimeUnit},
error::Result,
};
use super::iterator::{BufStreamingIterator, StreamingIterator};
use crate::array::{DictionaryArray, DictionaryKey, Offset};
use std::any::Any;
#[derive(Debug, PartialEq, Eq, Hash, Clone)]
pub struct SerializeOptions {
pub date32_format: String,
pub date64_format: String,
pub time32_format: String,
pub time64_format: String,
pub timestamp_format: String,
}
impl Default for SerializeOptions {
fn default() -> Self {
Self {
date32_format: "%F".to_string(),
date64_format: "%F".to_string(),
time32_format: "%T".to_string(),
time64_format: "%T".to_string(),
timestamp_format: "%FT%H:%M:%S.%9f".to_string(),
}
}
}
fn primitive_write<'a, T: NativeType + ToLexical>(
array: &'a PrimitiveArray<T>,
) -> Box<dyn StreamingIterator<Item = [u8]> + 'a> {
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
lexical_to_bytes_mut(*x, buf)
}
},
vec![],
))
}
macro_rules! dyn_primitive {
($ty:ty, $array:expr) => {{
let array = $array.as_any().downcast_ref().unwrap();
primitive_write::<$ty>(array)
}};
}
macro_rules! dyn_date {
($ty:ident, $fn:expr, $array:expr, $format:expr) => {{
let array = $array
.as_any()
.downcast_ref::<PrimitiveArray<$ty>>()
.unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
move |x, buf| {
if let Some(x) = x {
buf.extend_from_slice(($fn)(*x).format($format).to_string().as_bytes())
}
},
vec![],
))
}};
}
pub fn new_serializer<'a>(
array: &'a dyn Array,
options: &'a SerializeOptions,
) -> Result<Box<dyn StreamingIterator<Item = [u8]> + 'a>> {
Ok(match array.data_type() {
DataType::Boolean => {
let array = array.as_any().downcast_ref::<BooleanArray>().unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
if x {
buf.extend_from_slice(b"true");
} else {
buf.extend_from_slice(b"false");
}
}
},
vec![],
))
}
DataType::UInt8 => {
dyn_primitive!(u8, array)
}
DataType::UInt16 => {
dyn_primitive!(u16, array)
}
DataType::UInt32 => {
dyn_primitive!(u32, array)
}
DataType::UInt64 => {
dyn_primitive!(u64, array)
}
DataType::Int8 => {
dyn_primitive!(i8, array)
}
DataType::Int16 => {
dyn_primitive!(i16, array)
}
DataType::Int32 => {
dyn_primitive!(i32, array)
}
DataType::Date32 => {
dyn_date!(
i32,
temporal_conversions::date32_to_datetime,
array,
&options.date32_format
)
}
DataType::Time32(TimeUnit::Second) => {
dyn_date!(
i32,
temporal_conversions::time32s_to_time,
array,
&options.time32_format
)
}
DataType::Time32(TimeUnit::Millisecond) => {
dyn_date!(
i32,
temporal_conversions::time32ms_to_time,
array,
&options.time32_format
)
}
DataType::Int64 => {
dyn_primitive!(i64, array)
}
DataType::Date64 => {
dyn_date!(
i64,
temporal_conversions::date64_to_datetime,
array,
&options.date64_format
)
}
DataType::Time64(TimeUnit::Microsecond) => {
dyn_date!(
i64,
temporal_conversions::time64us_to_time,
array,
&options.time64_format
)
}
DataType::Time64(TimeUnit::Nanosecond) => {
dyn_date!(
i64,
temporal_conversions::time64ns_to_time,
array,
&options.time64_format
)
}
DataType::Timestamp(TimeUnit::Second, None) => {
dyn_date!(
i64,
temporal_conversions::timestamp_s_to_datetime,
array,
&options.timestamp_format
)
}
DataType::Timestamp(TimeUnit::Millisecond, None) => {
dyn_date!(
i64,
temporal_conversions::timestamp_ms_to_datetime,
array,
&options.timestamp_format
)
}
DataType::Timestamp(TimeUnit::Microsecond, None) => {
dyn_date!(
i64,
temporal_conversions::timestamp_us_to_datetime,
array,
&options.timestamp_format
)
}
DataType::Timestamp(TimeUnit::Nanosecond, None) => {
dyn_date!(
i64,
temporal_conversions::timestamp_ns_to_datetime,
array,
&options.timestamp_format
)
}
DataType::Float32 => {
dyn_primitive!(f32, array)
}
DataType::Float64 => {
dyn_primitive!(f64, array)
}
DataType::Utf8 => {
let array = array.as_any().downcast_ref::<Utf8Array<i32>>().unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
buf.extend_from_slice(x.as_bytes());
}
},
vec![],
))
}
DataType::LargeUtf8 => {
let array = array.as_any().downcast_ref::<Utf8Array<i64>>().unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
buf.extend_from_slice(x.as_bytes());
}
},
vec![],
))
}
DataType::Binary => {
let array = array.as_any().downcast_ref::<BinaryArray<i32>>().unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
buf.extend_from_slice(x);
}
},
vec![],
))
}
DataType::LargeBinary => {
let array = array.as_any().downcast_ref::<BinaryArray<i64>>().unwrap();
Box::new(BufStreamingIterator::new(
array.iter(),
|x, buf| {
if let Some(x) = x {
buf.extend_from_slice(x);
}
},
vec![],
))
}
DataType::Dictionary(keys_dt, values_dt) => match &**values_dt {
DataType::LargeUtf8 => match &**keys_dt {
DataType::UInt32 => serialize_utf8_dict::<u32, i64>(array.as_any()),
DataType::UInt64 => serialize_utf8_dict::<u64, i64>(array.as_any()),
_ => todo!(),
},
DataType::Utf8 => match &**keys_dt {
DataType::UInt32 => serialize_utf8_dict::<u32, i32>(array.as_any()),
DataType::UInt64 => serialize_utf8_dict::<u64, i32>(array.as_any()),
_ => todo!(),
},
_ => {
panic!("only dictionary with string values are supported by csv writer")
}
},
dt => panic!("data type: {} not supported by csv writer", dt),
})
}
fn serialize_utf8_dict<'a, K: DictionaryKey + Index, O: Offset>(
array: &'a dyn Any,
) -> Box<dyn StreamingIterator<Item = [u8]> + 'a> {
let array = array.downcast_ref::<DictionaryArray<K>>().unwrap();
let keys = array.keys();
let values = array
.values()
.as_any()
.downcast_ref::<Utf8Array<O>>()
.unwrap();
Box::new(BufStreamingIterator::new(
keys.iter(),
move |x, buf| {
if let Some(x) = x {
let i = Index::to_usize(x);
if !values.is_null(i) {
let val = values.value(i);
buf.extend_from_slice(val.as_bytes());
}
}
},
vec![],
))
}