1use std::io::Write;
18use std::sync::Arc;
19
20use crate::StructMode;
21use arrow_array::cast::AsArray;
22use arrow_array::types::*;
23use arrow_array::*;
24use arrow_buffer::{ArrowNativeType, NullBuffer, OffsetBuffer, ScalarBuffer};
25use arrow_cast::display::{ArrayFormatter, FormatOptions};
26use arrow_schema::{ArrowError, DataType, FieldRef};
27use half::f16;
28use lexical_core::FormattedSize;
29use serde_core::Serializer;
30
31#[derive(Debug, Clone, Default)]
33pub struct EncoderOptions {
34 explicit_nulls: bool,
36 struct_mode: StructMode,
38 encoder_factory: Option<Arc<dyn EncoderFactory>>,
40 date_format: Option<String>,
42 datetime_format: Option<String>,
44 timestamp_format: Option<String>,
46 timestamp_tz_format: Option<String>,
48 time_format: Option<String>,
50}
51
52impl EncoderOptions {
53 pub fn with_explicit_nulls(mut self, explicit_nulls: bool) -> Self {
55 self.explicit_nulls = explicit_nulls;
56 self
57 }
58
59 pub fn with_struct_mode(mut self, struct_mode: StructMode) -> Self {
61 self.struct_mode = struct_mode;
62 self
63 }
64
65 pub fn with_encoder_factory(mut self, encoder_factory: Arc<dyn EncoderFactory>) -> Self {
67 self.encoder_factory = Some(encoder_factory);
68 self
69 }
70
71 pub fn explicit_nulls(&self) -> bool {
73 self.explicit_nulls
74 }
75
76 pub fn struct_mode(&self) -> StructMode {
78 self.struct_mode
79 }
80
81 pub fn encoder_factory(&self) -> Option<&Arc<dyn EncoderFactory>> {
83 self.encoder_factory.as_ref()
84 }
85
86 pub fn with_date_format(mut self, format: String) -> Self {
88 self.date_format = Some(format);
89 self
90 }
91
92 pub fn date_format(&self) -> Option<&str> {
94 self.date_format.as_deref()
95 }
96
97 pub fn with_datetime_format(mut self, format: String) -> Self {
99 self.datetime_format = Some(format);
100 self
101 }
102
103 pub fn datetime_format(&self) -> Option<&str> {
105 self.datetime_format.as_deref()
106 }
107
108 pub fn with_time_format(mut self, format: String) -> Self {
110 self.time_format = Some(format);
111 self
112 }
113
114 pub fn time_format(&self) -> Option<&str> {
116 self.time_format.as_deref()
117 }
118
119 pub fn with_timestamp_format(mut self, format: String) -> Self {
121 self.timestamp_format = Some(format);
122 self
123 }
124
125 pub fn timestamp_format(&self) -> Option<&str> {
127 self.timestamp_format.as_deref()
128 }
129
130 pub fn with_timestamp_tz_format(mut self, tz_format: String) -> Self {
132 self.timestamp_tz_format = Some(tz_format);
133 self
134 }
135
136 pub fn timestamp_tz_format(&self) -> Option<&str> {
138 self.timestamp_tz_format.as_deref()
139 }
140}
141
142pub trait EncoderFactory: std::fmt::Debug + Send + Sync {
259 fn make_default_encoder<'a>(
267 &self,
268 _field: &'a FieldRef,
269 _array: &'a dyn Array,
270 _options: &'a EncoderOptions,
271 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
272 Ok(None)
273 }
274}
275
276pub struct NullableEncoder<'a> {
279 encoder: Box<dyn Encoder + 'a>,
280 nulls: Option<NullBuffer>,
281}
282
283impl<'a> NullableEncoder<'a> {
284 #[inline]
286 pub fn new(encoder: Box<dyn Encoder + 'a>, nulls: Option<NullBuffer>) -> Self {
287 Self { encoder, nulls }
288 }
289
290 #[inline]
292 pub fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
293 self.encoder.encode(idx, out)
294 }
295
296 #[inline]
298 pub fn is_null(&self, idx: usize) -> bool {
299 match self.nulls {
300 Some(ref nulls) => nulls.is_null(idx),
301 None => false,
302 }
303 }
304
305 #[inline]
307 pub fn has_nulls(&self) -> bool {
308 match self.nulls {
309 Some(ref nulls) => nulls.null_count() > 0,
310 None => false,
311 }
312 }
313}
314
315impl Encoder for NullableEncoder<'_> {
316 #[inline]
317 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
318 self.encoder.encode(idx, out)
319 }
320}
321
322pub trait Encoder {
326 fn encode(&mut self, idx: usize, out: &mut Vec<u8>);
330}
331
332pub fn make_encoder<'a>(
336 field: &'a FieldRef,
337 array: &'a dyn Array,
338 options: &'a EncoderOptions,
339) -> Result<NullableEncoder<'a>, ArrowError> {
340 macro_rules! primitive_helper {
341 ($t:ty) => {{
342 let array = array.as_primitive::<$t>();
343 let nulls = array.nulls().cloned();
344 NullableEncoder::new(Box::new(PrimitiveEncoder::new(array)), nulls)
345 }};
346 }
347
348 if let Some(factory) = options.encoder_factory()
349 && let Some(encoder) = factory.make_default_encoder(field, array, options)?
350 {
351 return Ok(encoder);
352 }
353
354 let nulls = array.nulls().cloned();
355 let encoder = downcast_integer! {
356 array.data_type() => (primitive_helper),
357 DataType::Float16 => primitive_helper!(Float16Type),
358 DataType::Float32 => primitive_helper!(Float32Type),
359 DataType::Float64 => primitive_helper!(Float64Type),
360 DataType::Boolean => {
361 let array = array.as_boolean();
362 NullableEncoder::new(Box::new(BooleanEncoder(array)), array.nulls().cloned())
363 }
364 DataType::Null => NullableEncoder::new(Box::new(NullEncoder), array.logical_nulls()),
365 DataType::Utf8 => {
366 let array = array.as_string::<i32>();
367 NullableEncoder::new(Box::new(StringEncoder(array)), array.nulls().cloned())
368 }
369 DataType::LargeUtf8 => {
370 let array = array.as_string::<i64>();
371 NullableEncoder::new(Box::new(StringEncoder(array)), array.nulls().cloned())
372 }
373 DataType::Utf8View => {
374 let array = array.as_string_view();
375 NullableEncoder::new(Box::new(StringViewEncoder(array)), array.nulls().cloned())
376 }
377 DataType::BinaryView => {
378 let array = array.as_binary_view();
379 NullableEncoder::new(Box::new(BinaryViewEncoder(array)), array.nulls().cloned())
380 }
381 DataType::List(_) => {
382 let array = array.as_list::<i32>();
383 NullableEncoder::new(Box::new(ListLikeEncoder::try_new(field, array, options)?), array.nulls().cloned())
384 }
385 DataType::LargeList(_) => {
386 let array = array.as_list::<i64>();
387 NullableEncoder::new(Box::new(ListLikeEncoder::try_new(field, array, options)?), array.nulls().cloned())
388 }
389 DataType::ListView(_) => {
390 let array = array.as_list_view::<i32>();
391 NullableEncoder::new(Box::new(ListLikeEncoder::try_new(field, array, options)?), array.nulls().cloned())
392 }
393 DataType::LargeListView(_) => {
394 let array = array.as_list_view::<i64>();
395 NullableEncoder::new(Box::new(ListLikeEncoder::try_new(field, array, options)?), array.nulls().cloned())
396 }
397 DataType::FixedSizeList(_, _) => {
398 let array = array.as_fixed_size_list();
399 NullableEncoder::new(Box::new(ListLikeEncoder::try_new(field, array, options)?), array.nulls().cloned())
400 }
401
402 DataType::Dictionary(_, _) => downcast_dictionary_array! {
403 array => {
404 NullableEncoder::new(Box::new(DictionaryEncoder::try_new(field, array, options)?), array.nulls().cloned())
405 },
406 _ => unreachable!()
407 }
408
409 DataType::RunEndEncoded(_, _) => downcast_run_array! {
410 array => {
411 NullableEncoder::new(
412 Box::new(RunEndEncodedEncoder::try_new(field, array, options)?),
413 array.logical_nulls(),
414 )
415 },
416 _ => unreachable!()
417 }
418
419 DataType::Map(_, _) => {
420 let array = array.as_map();
421 NullableEncoder::new(Box::new(MapEncoder::try_new(field, array, options)?), array.nulls().cloned())
422 }
423
424 DataType::FixedSizeBinary(_) => {
425 let array = array.as_fixed_size_binary();
426 NullableEncoder::new(Box::new(BinaryEncoder::new(array)) as _, array.nulls().cloned())
427 }
428
429 DataType::Binary => {
430 let array: &BinaryArray = array.as_binary();
431 NullableEncoder::new(Box::new(BinaryEncoder::new(array)), array.nulls().cloned())
432 }
433
434 DataType::LargeBinary => {
435 let array: &LargeBinaryArray = array.as_binary();
436 NullableEncoder::new(Box::new(BinaryEncoder::new(array)), array.nulls().cloned())
437 }
438
439 DataType::Struct(fields) => {
440 let array = array.as_struct();
441 let encoders = fields.iter().zip(array.columns()).map(|(field, array)| {
442 let encoder = make_encoder(field, array, options)?;
443
444 let mut field_name = Vec::with_capacity(field.name().len() + 3);
446 encode_string(field.name(), &mut field_name);
447 field_name.push(b':');
448
449 Ok(FieldEncoder {
450 field_name,
451 encoder,
452 })
453 }).collect::<Result<Vec<_>, ArrowError>>()?;
454
455 let encoder = StructArrayEncoder{
456 encoders,
457 explicit_nulls: options.explicit_nulls(),
458 struct_mode: options.struct_mode(),
459 };
460 let nulls = array.nulls().cloned();
461 NullableEncoder::new(Box::new(encoder) as Box<dyn Encoder + 'a>, nulls)
462 }
463 DataType::Decimal32(_, _) | DataType::Decimal64(_, _) | DataType::Decimal128(_, _) | DataType::Decimal256(_, _) => {
464 let options = FormatOptions::new().with_display_error(true);
465 let formatter = JsonArrayFormatter::new(ArrayFormatter::try_new(array, &options)?);
466 NullableEncoder::new(Box::new(RawArrayFormatter(formatter)) as Box<dyn Encoder + 'a>, nulls)
467 }
468 d => match d.is_temporal() {
469 true => {
470 let fops = FormatOptions::new().with_display_error(true)
475 .with_date_format(options.date_format.as_deref())
476 .with_datetime_format(options.datetime_format.as_deref())
477 .with_timestamp_format(options.timestamp_format.as_deref())
478 .with_timestamp_tz_format(options.timestamp_tz_format.as_deref())
479 .with_time_format(options.time_format.as_deref());
480
481 let formatter = ArrayFormatter::try_new(array, &fops)?;
482 let formatter = JsonArrayFormatter::new(formatter);
483 NullableEncoder::new(Box::new(formatter) as Box<dyn Encoder + 'a>, nulls)
484 }
485 false => return Err(ArrowError::JsonError(format!(
486 "Unsupported data type for JSON encoding: {d:?}",
487 )))
488 }
489 };
490
491 Ok(encoder)
492}
493
494fn encode_string(s: &str, out: &mut Vec<u8>) {
495 let mut serializer = serde_json::Serializer::new(out);
496 serializer.serialize_str(s).unwrap();
497}
498
499fn encode_binary(bytes: &[u8], out: &mut Vec<u8>) {
500 out.push(b'"');
501 for byte in bytes {
502 write!(out, "{byte:02x}").unwrap();
503 }
504 out.push(b'"');
505}
506
507struct FieldEncoder<'a> {
508 field_name: Vec<u8>,
509 encoder: NullableEncoder<'a>,
510}
511
512impl FieldEncoder<'_> {
513 #[inline]
514 fn is_null(&self, idx: usize) -> bool {
515 self.encoder.is_null(idx)
516 }
517}
518
519struct StructArrayEncoder<'a> {
520 encoders: Vec<FieldEncoder<'a>>,
521 explicit_nulls: bool,
522 struct_mode: StructMode,
523}
524
525impl Encoder for StructArrayEncoder<'_> {
526 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
527 match self.struct_mode {
528 StructMode::ObjectOnly => out.push(b'{'),
529 StructMode::ListOnly => out.push(b'['),
530 }
531 let mut is_first = true;
532 let drop_nulls = (self.struct_mode == StructMode::ObjectOnly) && !self.explicit_nulls;
534
535 for field_encoder in &mut self.encoders {
536 let is_null = field_encoder.is_null(idx);
537 if is_null && drop_nulls {
538 continue;
539 }
540
541 if !is_first {
542 out.push(b',');
543 }
544 is_first = false;
545
546 if self.struct_mode == StructMode::ObjectOnly {
547 out.extend_from_slice(&field_encoder.field_name);
548 }
549
550 if is_null {
551 out.extend_from_slice(b"null");
552 } else {
553 field_encoder.encoder.encode(idx, out);
554 }
555 }
556 match self.struct_mode {
557 StructMode::ObjectOnly => out.push(b'}'),
558 StructMode::ListOnly => out.push(b']'),
559 }
560 }
561}
562
563trait PrimitiveEncode: ArrowNativeType {
564 type Buffer;
565
566 fn init_buffer() -> Self::Buffer;
568
569 fn encode(self, buf: &mut Self::Buffer) -> &[u8];
573}
574
575macro_rules! integer_encode {
576 ($($t:ty),*) => {
577 $(
578 impl PrimitiveEncode for $t {
579 type Buffer = [u8; Self::FORMATTED_SIZE];
580
581 fn init_buffer() -> Self::Buffer {
582 [0; Self::FORMATTED_SIZE]
583 }
584
585 fn encode(self, buf: &mut Self::Buffer) -> &[u8] {
586 lexical_core::write(self, buf)
587 }
588 }
589 )*
590 };
591}
592integer_encode!(i8, i16, i32, i64, u8, u16, u32, u64);
593
594macro_rules! float_encode {
595 ($($t:ty),*) => {
596 $(
597 impl PrimitiveEncode for $t {
598 type Buffer = [u8; Self::FORMATTED_SIZE];
599
600 fn init_buffer() -> Self::Buffer {
601 [0; Self::FORMATTED_SIZE]
602 }
603
604 fn encode(self, buf: &mut Self::Buffer) -> &[u8] {
605 if self.is_infinite() || self.is_nan() {
606 b"null"
607 } else {
608 lexical_core::write(self, buf)
609 }
610 }
611 }
612 )*
613 };
614}
615float_encode!(f32, f64);
616
617impl PrimitiveEncode for f16 {
618 type Buffer = <f32 as PrimitiveEncode>::Buffer;
619
620 fn init_buffer() -> Self::Buffer {
621 f32::init_buffer()
622 }
623
624 fn encode(self, buf: &mut Self::Buffer) -> &[u8] {
625 self.to_f32().encode(buf)
626 }
627}
628
629struct PrimitiveEncoder<N: PrimitiveEncode> {
630 values: ScalarBuffer<N>,
631 buffer: N::Buffer,
632}
633
634impl<N: PrimitiveEncode> PrimitiveEncoder<N> {
635 fn new<P: ArrowPrimitiveType<Native = N>>(array: &PrimitiveArray<P>) -> Self {
636 Self {
637 values: array.values().clone(),
638 buffer: N::init_buffer(),
639 }
640 }
641}
642
643impl<N: PrimitiveEncode> Encoder for PrimitiveEncoder<N> {
644 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
645 out.extend_from_slice(self.values[idx].encode(&mut self.buffer));
646 }
647}
648
649struct BooleanEncoder<'a>(&'a BooleanArray);
650
651impl Encoder for BooleanEncoder<'_> {
652 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
653 match self.0.value(idx) {
654 true => out.extend_from_slice(b"true"),
655 false => out.extend_from_slice(b"false"),
656 }
657 }
658}
659
660struct StringEncoder<'a, O: OffsetSizeTrait>(&'a GenericStringArray<O>);
661
662impl<O: OffsetSizeTrait> Encoder for StringEncoder<'_, O> {
663 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
664 encode_string(self.0.value(idx), out);
665 }
666}
667
668struct StringViewEncoder<'a>(&'a StringViewArray);
669
670impl Encoder for StringViewEncoder<'_> {
671 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
672 encode_string(self.0.value(idx), out);
673 }
674}
675
676struct BinaryViewEncoder<'a>(&'a BinaryViewArray);
677
678impl Encoder for BinaryViewEncoder<'_> {
679 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
680 encode_binary(self.0.value(idx), out);
681 }
682}
683
684struct ListLikeEncoder<'a, L: ListLikeArray> {
685 list_array: &'a L,
686 encoder: NullableEncoder<'a>,
687}
688
689impl<'a, L: ListLikeArray> ListLikeEncoder<'a, L> {
690 fn try_new(
691 field: &'a FieldRef,
692 array: &'a L,
693 options: &'a EncoderOptions,
694 ) -> Result<Self, ArrowError> {
695 let encoder = make_encoder(field, array.values().as_ref(), options)?;
696 Ok(Self {
697 list_array: array,
698 encoder,
699 })
700 }
701}
702
703impl<L: ListLikeArray> Encoder for ListLikeEncoder<'_, L> {
704 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
705 let range = self.list_array.element_range(idx);
706 let start = range.start;
707 let end = range.end;
708 out.push(b'[');
709 if self.encoder.has_nulls() {
710 for idx in start..end {
711 if idx != start {
712 out.push(b',')
713 }
714 if self.encoder.is_null(idx) {
715 out.extend_from_slice(b"null");
716 } else {
717 self.encoder.encode(idx, out);
718 }
719 }
720 } else {
721 for idx in start..end {
722 if idx != start {
723 out.push(b',')
724 }
725 self.encoder.encode(idx, out);
726 }
727 }
728 out.push(b']');
729 }
730}
731
732struct DictionaryEncoder<'a, K: ArrowDictionaryKeyType> {
733 keys: ScalarBuffer<K::Native>,
734 encoder: NullableEncoder<'a>,
735}
736
737impl<'a, K: ArrowDictionaryKeyType> DictionaryEncoder<'a, K> {
738 fn try_new(
739 field: &'a FieldRef,
740 array: &'a DictionaryArray<K>,
741 options: &'a EncoderOptions,
742 ) -> Result<Self, ArrowError> {
743 let encoder = make_encoder(field, array.values().as_ref(), options)?;
744
745 Ok(Self {
746 keys: array.keys().values().clone(),
747 encoder,
748 })
749 }
750}
751
752impl<K: ArrowDictionaryKeyType> Encoder for DictionaryEncoder<'_, K> {
753 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
754 self.encoder.encode(self.keys[idx].as_usize(), out)
755 }
756}
757
758struct RunEndEncodedEncoder<'a, R: RunEndIndexType> {
759 run_array: &'a RunArray<R>,
760 encoder: NullableEncoder<'a>,
761}
762
763impl<'a, R: RunEndIndexType> RunEndEncodedEncoder<'a, R> {
764 fn try_new(
765 field: &'a FieldRef,
766 array: &'a RunArray<R>,
767 options: &'a EncoderOptions,
768 ) -> Result<Self, ArrowError> {
769 let encoder = make_encoder(field, array.values().as_ref(), options)?;
770 Ok(Self {
771 run_array: array,
772 encoder,
773 })
774 }
775}
776
777impl<R: RunEndIndexType> Encoder for RunEndEncodedEncoder<'_, R> {
778 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
779 let physical_idx = self.run_array.get_physical_index(idx);
780 self.encoder.encode(physical_idx, out)
781 }
782}
783
784struct JsonArrayFormatter<'a> {
786 formatter: ArrayFormatter<'a>,
787}
788
789impl<'a> JsonArrayFormatter<'a> {
790 fn new(formatter: ArrayFormatter<'a>) -> Self {
791 Self { formatter }
792 }
793}
794
795impl Encoder for JsonArrayFormatter<'_> {
796 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
797 out.push(b'"');
798 let _ = write!(out, "{}", self.formatter.value(idx));
801 out.push(b'"')
802 }
803}
804
805struct RawArrayFormatter<'a>(JsonArrayFormatter<'a>);
807
808impl Encoder for RawArrayFormatter<'_> {
809 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
810 let _ = write!(out, "{}", self.0.formatter.value(idx));
811 }
812}
813
814struct NullEncoder;
815
816impl Encoder for NullEncoder {
817 fn encode(&mut self, _idx: usize, _out: &mut Vec<u8>) {
818 unreachable!()
819 }
820}
821
822struct MapEncoder<'a> {
823 offsets: OffsetBuffer<i32>,
824 keys: NullableEncoder<'a>,
825 values: NullableEncoder<'a>,
826 explicit_nulls: bool,
827}
828
829impl<'a> MapEncoder<'a> {
830 fn try_new(
831 field: &'a FieldRef,
832 array: &'a MapArray,
833 options: &'a EncoderOptions,
834 ) -> Result<Self, ArrowError> {
835 let values = array.values();
836 let keys = array.keys();
837
838 if !matches!(
839 keys.data_type(),
840 DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View
841 ) {
842 return Err(ArrowError::JsonError(format!(
843 "Only UTF8 keys supported by JSON MapArray Writer: got {:?}",
844 keys.data_type()
845 )));
846 }
847
848 let keys = make_encoder(field, keys, options)?;
849 let values = make_encoder(field, values, options)?;
850
851 if keys.has_nulls() {
853 return Err(ArrowError::InvalidArgumentError(
854 "Encountered nulls in MapArray keys".to_string(),
855 ));
856 }
857
858 if array.entries().nulls().is_some_and(|x| x.null_count() != 0) {
859 return Err(ArrowError::InvalidArgumentError(
860 "Encountered nulls in MapArray entries".to_string(),
861 ));
862 }
863
864 Ok(Self {
865 offsets: array.offsets().clone(),
866 keys,
867 values,
868 explicit_nulls: options.explicit_nulls(),
869 })
870 }
871}
872
873impl Encoder for MapEncoder<'_> {
874 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
875 let end = self.offsets[idx + 1].as_usize();
876 let start = self.offsets[idx].as_usize();
877
878 let mut is_first = true;
879
880 out.push(b'{');
881
882 for idx in start..end {
883 let is_null = self.values.is_null(idx);
884 if is_null && !self.explicit_nulls {
885 continue;
886 }
887
888 if !is_first {
889 out.push(b',');
890 }
891 is_first = false;
892
893 self.keys.encode(idx, out);
894 out.push(b':');
895
896 if is_null {
897 out.extend_from_slice(b"null");
898 } else {
899 self.values.encode(idx, out);
900 }
901 }
902 out.push(b'}');
903 }
904}
905
906struct BinaryEncoder<B>(B);
909
910impl<'a, B> BinaryEncoder<B>
911where
912 B: ArrayAccessor<Item = &'a [u8]>,
913{
914 fn new(array: B) -> Self {
915 Self(array)
916 }
917}
918
919impl<'a, B> Encoder for BinaryEncoder<B>
920where
921 B: ArrayAccessor<Item = &'a [u8]>,
922{
923 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
924 out.push(b'"');
925 for byte in self.0.value(idx) {
926 write!(out, "{byte:02x}").unwrap();
928 }
929 out.push(b'"');
930 }
931}