1mod encoder;
108
109use std::{fmt::Debug, io::Write, sync::Arc};
110
111use crate::StructMode;
112use arrow_array::*;
113use arrow_schema::*;
114
115pub use encoder::{Encoder, EncoderFactory, EncoderOptions, NullableEncoder, make_encoder};
116
117pub trait JsonFormat: Debug + Default {
120 #[inline]
121 fn start_stream<W: Write>(&self, _writer: &mut W) -> Result<(), ArrowError> {
123 Ok(())
124 }
125
126 #[inline]
127 fn start_row<W: Write>(&self, _writer: &mut W, _is_first_row: bool) -> Result<(), ArrowError> {
129 Ok(())
130 }
131
132 #[inline]
133 fn end_row<W: Write>(&self, _writer: &mut W) -> Result<(), ArrowError> {
135 Ok(())
136 }
137
138 fn end_stream<W: Write>(&self, _writer: &mut W) -> Result<(), ArrowError> {
140 Ok(())
141 }
142}
143
144#[derive(Debug, Default)]
154pub struct LineDelimited {}
155
156impl JsonFormat for LineDelimited {
157 fn end_row<W: Write>(&self, writer: &mut W) -> Result<(), ArrowError> {
158 writer.write_all(b"\n")?;
159 Ok(())
160 }
161}
162
163#[derive(Debug, Default)]
171pub struct JsonArray {}
172
173impl JsonFormat for JsonArray {
174 fn start_stream<W: Write>(&self, writer: &mut W) -> Result<(), ArrowError> {
175 writer.write_all(b"[")?;
176 Ok(())
177 }
178
179 fn start_row<W: Write>(&self, writer: &mut W, is_first_row: bool) -> Result<(), ArrowError> {
180 if !is_first_row {
181 writer.write_all(b",")?;
182 }
183 Ok(())
184 }
185
186 fn end_stream<W: Write>(&self, writer: &mut W) -> Result<(), ArrowError> {
187 writer.write_all(b"]")?;
188 Ok(())
189 }
190}
191
192pub type LineDelimitedWriter<W> = Writer<W, LineDelimited>;
194
195pub type ArrayWriter<W> = Writer<W, JsonArray>;
197
198#[derive(Debug, Clone, Default)]
200pub struct WriterBuilder(EncoderOptions);
201
202impl WriterBuilder {
203 pub fn new() -> Self {
223 Self::default()
224 }
225
226 pub fn explicit_nulls(&self) -> bool {
228 self.0.explicit_nulls()
229 }
230
231 pub fn with_explicit_nulls(mut self, explicit_nulls: bool) -> Self {
254 self.0 = self.0.with_explicit_nulls(explicit_nulls);
255 self
256 }
257
258 pub fn struct_mode(&self) -> StructMode {
260 self.0.struct_mode()
261 }
262
263 pub fn with_struct_mode(mut self, struct_mode: StructMode) -> Self {
269 self.0 = self.0.with_struct_mode(struct_mode);
270 self
271 }
272
273 pub fn with_encoder_factory(mut self, factory: Arc<dyn EncoderFactory>) -> Self {
278 self.0 = self.0.with_encoder_factory(factory);
279 self
280 }
281
282 pub fn with_date_format(mut self, format: String) -> Self {
284 self.0 = self.0.with_date_format(format);
285 self
286 }
287
288 pub fn with_datetime_format(mut self, format: String) -> Self {
290 self.0 = self.0.with_datetime_format(format);
291 self
292 }
293
294 pub fn with_time_format(mut self, format: String) -> Self {
296 self.0 = self.0.with_time_format(format);
297 self
298 }
299
300 pub fn with_timestamp_format(mut self, format: String) -> Self {
302 self.0 = self.0.with_timestamp_format(format);
303 self
304 }
305
306 pub fn with_timestamp_tz_format(mut self, tz_format: String) -> Self {
308 self.0 = self.0.with_timestamp_tz_format(tz_format);
309 self
310 }
311
312 pub fn build<W, F>(self, writer: W) -> Writer<W, F>
314 where
315 W: Write,
316 F: JsonFormat,
317 {
318 Writer {
319 writer,
320 started: false,
321 finished: false,
322 format: F::default(),
323 options: self.0,
324 }
325 }
326}
327
328#[derive(Debug)]
339pub struct Writer<W, F>
340where
341 W: Write,
342 F: JsonFormat,
343{
344 writer: W,
346
347 started: bool,
349
350 finished: bool,
352
353 format: F,
355
356 options: EncoderOptions,
358}
359
360impl<W, F> Writer<W, F>
361where
362 W: Write,
363 F: JsonFormat,
364{
365 pub fn new(writer: W) -> Self {
367 Self {
368 writer,
369 started: false,
370 finished: false,
371 format: F::default(),
372 options: EncoderOptions::default(),
373 }
374 }
375
376 pub fn write(&mut self, batch: &RecordBatch) -> Result<(), ArrowError> {
378 if batch.num_rows() == 0 {
379 return Ok(());
380 }
381
382 let mut buffer = Vec::with_capacity(16 * 1024);
385
386 let mut is_first_row = !self.started;
387 if !self.started {
388 self.format.start_stream(&mut buffer)?;
389 self.started = true;
390 }
391
392 let array = StructArray::from(batch.clone());
393 let field = Arc::new(Field::new_struct(
394 "",
395 batch.schema().fields().clone(),
396 false,
397 ));
398
399 let mut encoder = make_encoder(&field, &array, &self.options)?;
400
401 assert!(!encoder.has_nulls(), "root cannot be nullable");
403 for idx in 0..batch.num_rows() {
404 self.format.start_row(&mut buffer, is_first_row)?;
405 is_first_row = false;
406
407 encoder.encode(idx, &mut buffer);
408 if buffer.len() > 8 * 1024 {
409 self.writer.write_all(&buffer)?;
410 buffer.clear();
411 }
412 self.format.end_row(&mut buffer)?;
413 }
414
415 if !buffer.is_empty() {
416 self.writer.write_all(&buffer)?;
417 }
418
419 Ok(())
420 }
421
422 pub fn write_batches(&mut self, batches: &[&RecordBatch]) -> Result<(), ArrowError> {
424 for b in batches {
425 self.write(b)?;
426 }
427 Ok(())
428 }
429
430 pub fn finish(&mut self) -> Result<(), ArrowError> {
434 if !self.started {
435 self.format.start_stream(&mut self.writer)?;
436 self.started = true;
437 }
438 if !self.finished {
439 self.format.end_stream(&mut self.writer)?;
440 self.finished = true;
441 }
442
443 Ok(())
444 }
445
446 pub fn get_ref(&self) -> &W {
448 &self.writer
449 }
450
451 pub fn get_mut(&mut self) -> &mut W {
456 &mut self.writer
457 }
458
459 pub fn into_inner(self) -> W {
461 self.writer
462 }
463}
464
465impl<W, F> RecordBatchWriter for Writer<W, F>
466where
467 W: Write,
468 F: JsonFormat,
469{
470 fn write(&mut self, batch: &RecordBatch) -> Result<(), ArrowError> {
471 self.write(batch)
472 }
473
474 fn close(mut self) -> Result<(), ArrowError> {
475 self.finish()
476 }
477}
478
479#[cfg(test)]
480mod tests {
481 use core::str;
482 use std::collections::HashMap;
483 use std::fs::{File, read_to_string};
484 use std::io::{BufReader, Seek};
485 use std::sync::Arc;
486
487 use arrow_array::cast::AsArray;
488 use serde_json::{Value, json};
489
490 use super::LineDelimited;
491 use super::{Encoder, WriterBuilder};
492 use arrow_array::builder::*;
493 use arrow_array::types::*;
494 use arrow_buffer::{Buffer, NullBuffer, OffsetBuffer, ScalarBuffer, i256};
495
496 use crate::reader::*;
497
498 use super::*;
499
500 fn assert_json_eq(input: &[u8], expected: &str) {
502 let expected: Vec<Option<Value>> = expected
503 .split('\n')
504 .map(|s| (!s.is_empty()).then(|| serde_json::from_str(s).unwrap()))
505 .collect();
506
507 let actual: Vec<Option<Value>> = input
508 .split(|b| *b == b'\n')
509 .map(|s| (!s.is_empty()).then(|| serde_json::from_slice(s).unwrap()))
510 .collect();
511
512 assert_eq!(actual, expected);
513 }
514
515 #[test]
516 fn write_simple_rows() {
517 let schema = Schema::new(vec![
518 Field::new("c1", DataType::Int32, true),
519 Field::new("c2", DataType::Utf8, true),
520 ]);
521
522 let a = Int32Array::from(vec![Some(1), Some(2), Some(3), None, Some(5)]);
523 let b = StringArray::from(vec![Some("a"), Some("b"), Some("c"), Some("d"), None]);
524
525 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(a), Arc::new(b)]).unwrap();
526
527 let mut buf = Vec::new();
528 {
529 let mut writer = LineDelimitedWriter::new(&mut buf);
530 writer.write_batches(&[&batch]).unwrap();
531 }
532
533 assert_json_eq(
534 &buf,
535 r#"{"c1":1,"c2":"a"}
536{"c1":2,"c2":"b"}
537{"c1":3,"c2":"c"}
538{"c2":"d"}
539{"c1":5}
540"#,
541 );
542 }
543
544 #[test]
545 fn write_large_utf8_and_utf8_view() {
546 let schema = Schema::new(vec![
547 Field::new("c1", DataType::Utf8, true),
548 Field::new("c2", DataType::LargeUtf8, true),
549 Field::new("c3", DataType::Utf8View, true),
550 ]);
551
552 let a = StringArray::from(vec![Some("a"), None, Some("c"), Some("d"), None]);
553 let b = LargeStringArray::from(vec![Some("a"), Some("b"), None, Some("d"), None]);
554 let c = StringViewArray::from(vec![Some("a"), Some("b"), None, Some("d"), None]);
555
556 let batch = RecordBatch::try_new(
557 Arc::new(schema),
558 vec![Arc::new(a), Arc::new(b), Arc::new(c)],
559 )
560 .unwrap();
561
562 let mut buf = Vec::new();
563 {
564 let mut writer = LineDelimitedWriter::new(&mut buf);
565 writer.write_batches(&[&batch]).unwrap();
566 }
567
568 assert_json_eq(
569 &buf,
570 r#"{"c1":"a","c2":"a","c3":"a"}
571{"c2":"b","c3":"b"}
572{"c1":"c"}
573{"c1":"d","c2":"d","c3":"d"}
574{}
575"#,
576 );
577 }
578
579 #[test]
580 fn write_dictionary() {
581 let schema = Schema::new(vec![
582 Field::new_dictionary("c1", DataType::Int32, DataType::Utf8, true),
583 Field::new_dictionary("c2", DataType::Int8, DataType::Utf8, true),
584 ]);
585
586 let a: DictionaryArray<Int32Type> = vec![
587 Some("cupcakes"),
588 Some("foo"),
589 Some("foo"),
590 None,
591 Some("cupcakes"),
592 ]
593 .into_iter()
594 .collect();
595 let b: DictionaryArray<Int8Type> =
596 vec![Some("sdsd"), Some("sdsd"), None, Some("sd"), Some("sdsd")]
597 .into_iter()
598 .collect();
599
600 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(a), Arc::new(b)]).unwrap();
601
602 let mut buf = Vec::new();
603 {
604 let mut writer = LineDelimitedWriter::new(&mut buf);
605 writer.write_batches(&[&batch]).unwrap();
606 }
607
608 assert_json_eq(
609 &buf,
610 r#"{"c1":"cupcakes","c2":"sdsd"}
611{"c1":"foo","c2":"sdsd"}
612{"c1":"foo"}
613{"c2":"sd"}
614{"c1":"cupcakes","c2":"sdsd"}
615"#,
616 );
617 }
618
619 #[test]
620 fn write_list_of_dictionary() {
621 let dict_field = Arc::new(Field::new_dictionary(
622 "item",
623 DataType::Int32,
624 DataType::Utf8,
625 true,
626 ));
627 let schema = Schema::new(vec![Field::new_large_list("l", dict_field.clone(), true)]);
628
629 let dict_array: DictionaryArray<Int32Type> =
630 vec![Some("a"), Some("b"), Some("c"), Some("a"), None, Some("c")]
631 .into_iter()
632 .collect();
633 let list_array = LargeListArray::try_new(
634 dict_field,
635 OffsetBuffer::from_lengths([3_usize, 2, 0, 1]),
636 Arc::new(dict_array),
637 Some(NullBuffer::from_iter([true, true, false, true])),
638 )
639 .unwrap();
640
641 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(list_array)]).unwrap();
642
643 let mut buf = Vec::new();
644 {
645 let mut writer = LineDelimitedWriter::new(&mut buf);
646 writer.write_batches(&[&batch]).unwrap();
647 }
648
649 assert_json_eq(
650 &buf,
651 r#"{"l":["a","b","c"]}
652{"l":["a",null]}
653{}
654{"l":["c"]}
655"#,
656 );
657 }
658
659 #[test]
660 fn write_list_of_dictionary_large_values() {
661 let dict_field = Arc::new(Field::new_dictionary(
662 "item",
663 DataType::Int32,
664 DataType::LargeUtf8,
665 true,
666 ));
667 let schema = Schema::new(vec![Field::new_large_list("l", dict_field.clone(), true)]);
668
669 let keys = PrimitiveArray::<Int32Type>::from(vec![
670 Some(0),
671 Some(1),
672 Some(2),
673 Some(0),
674 None,
675 Some(2),
676 ]);
677 let values = LargeStringArray::from(vec!["a", "b", "c"]);
678 let dict_array = DictionaryArray::try_new(keys, Arc::new(values)).unwrap();
679
680 let list_array = LargeListArray::try_new(
681 dict_field,
682 OffsetBuffer::from_lengths([3_usize, 2, 0, 1]),
683 Arc::new(dict_array),
684 Some(NullBuffer::from_iter([true, true, false, true])),
685 )
686 .unwrap();
687
688 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(list_array)]).unwrap();
689
690 let mut buf = Vec::new();
691 {
692 let mut writer = LineDelimitedWriter::new(&mut buf);
693 writer.write_batches(&[&batch]).unwrap();
694 }
695
696 assert_json_eq(
697 &buf,
698 r#"{"l":["a","b","c"]}
699{"l":["a",null]}
700{}
701{"l":["c"]}
702"#,
703 );
704 }
705
706 #[test]
707 fn write_timestamps() {
708 let ts_string = "2018-11-13T17:11:10.011375885995";
709 let ts_nanos = ts_string
710 .parse::<chrono::NaiveDateTime>()
711 .unwrap()
712 .and_utc()
713 .timestamp_nanos_opt()
714 .unwrap();
715 let ts_micros = ts_nanos / 1000;
716 let ts_millis = ts_micros / 1000;
717 let ts_secs = ts_millis / 1000;
718
719 let arr_nanos = TimestampNanosecondArray::from(vec![Some(ts_nanos), None]);
720 let arr_micros = TimestampMicrosecondArray::from(vec![Some(ts_micros), None]);
721 let arr_millis = TimestampMillisecondArray::from(vec![Some(ts_millis), None]);
722 let arr_secs = TimestampSecondArray::from(vec![Some(ts_secs), None]);
723 let arr_names = StringArray::from(vec![Some("a"), Some("b")]);
724
725 let schema = Schema::new(vec![
726 Field::new("nanos", arr_nanos.data_type().clone(), true),
727 Field::new("micros", arr_micros.data_type().clone(), true),
728 Field::new("millis", arr_millis.data_type().clone(), true),
729 Field::new("secs", arr_secs.data_type().clone(), true),
730 Field::new("name", arr_names.data_type().clone(), true),
731 ]);
732 let schema = Arc::new(schema);
733
734 let batch = RecordBatch::try_new(
735 schema,
736 vec![
737 Arc::new(arr_nanos),
738 Arc::new(arr_micros),
739 Arc::new(arr_millis),
740 Arc::new(arr_secs),
741 Arc::new(arr_names),
742 ],
743 )
744 .unwrap();
745
746 let mut buf = Vec::new();
747 {
748 let mut writer = LineDelimitedWriter::new(&mut buf);
749 writer.write_batches(&[&batch]).unwrap();
750 }
751
752 assert_json_eq(
753 &buf,
754 r#"{"micros":"2018-11-13T17:11:10.011375","millis":"2018-11-13T17:11:10.011","name":"a","nanos":"2018-11-13T17:11:10.011375885","secs":"2018-11-13T17:11:10"}
755{"name":"b"}
756"#,
757 );
758
759 let mut buf = Vec::new();
760 {
761 let mut writer = WriterBuilder::new()
762 .with_timestamp_format("%m-%d-%Y".to_string())
763 .build::<_, LineDelimited>(&mut buf);
764 writer.write_batches(&[&batch]).unwrap();
765 }
766
767 assert_json_eq(
768 &buf,
769 r#"{"nanos":"11-13-2018","micros":"11-13-2018","millis":"11-13-2018","secs":"11-13-2018","name":"a"}
770{"name":"b"}
771"#,
772 );
773 }
774
775 #[test]
776 fn write_timestamps_with_tz() {
777 let ts_string = "2018-11-13T17:11:10.011375885995";
778 let ts_nanos = ts_string
779 .parse::<chrono::NaiveDateTime>()
780 .unwrap()
781 .and_utc()
782 .timestamp_nanos_opt()
783 .unwrap();
784 let ts_micros = ts_nanos / 1000;
785 let ts_millis = ts_micros / 1000;
786 let ts_secs = ts_millis / 1000;
787
788 let arr_nanos = TimestampNanosecondArray::from(vec![Some(ts_nanos), None]);
789 let arr_micros = TimestampMicrosecondArray::from(vec![Some(ts_micros), None]);
790 let arr_millis = TimestampMillisecondArray::from(vec![Some(ts_millis), None]);
791 let arr_secs = TimestampSecondArray::from(vec![Some(ts_secs), None]);
792 let arr_names = StringArray::from(vec![Some("a"), Some("b")]);
793
794 let tz = "+00:00";
795
796 let arr_nanos = arr_nanos.with_timezone(tz);
797 let arr_micros = arr_micros.with_timezone(tz);
798 let arr_millis = arr_millis.with_timezone(tz);
799 let arr_secs = arr_secs.with_timezone(tz);
800
801 let schema = Schema::new(vec![
802 Field::new("nanos", arr_nanos.data_type().clone(), true),
803 Field::new("micros", arr_micros.data_type().clone(), true),
804 Field::new("millis", arr_millis.data_type().clone(), true),
805 Field::new("secs", arr_secs.data_type().clone(), true),
806 Field::new("name", arr_names.data_type().clone(), true),
807 ]);
808 let schema = Arc::new(schema);
809
810 let batch = RecordBatch::try_new(
811 schema,
812 vec![
813 Arc::new(arr_nanos),
814 Arc::new(arr_micros),
815 Arc::new(arr_millis),
816 Arc::new(arr_secs),
817 Arc::new(arr_names),
818 ],
819 )
820 .unwrap();
821
822 let mut buf = Vec::new();
823 {
824 let mut writer = LineDelimitedWriter::new(&mut buf);
825 writer.write_batches(&[&batch]).unwrap();
826 }
827
828 assert_json_eq(
829 &buf,
830 r#"{"micros":"2018-11-13T17:11:10.011375Z","millis":"2018-11-13T17:11:10.011Z","name":"a","nanos":"2018-11-13T17:11:10.011375885Z","secs":"2018-11-13T17:11:10Z"}
831{"name":"b"}
832"#,
833 );
834
835 let mut buf = Vec::new();
836 {
837 let mut writer = WriterBuilder::new()
838 .with_timestamp_tz_format("%m-%d-%Y %Z".to_string())
839 .build::<_, LineDelimited>(&mut buf);
840 writer.write_batches(&[&batch]).unwrap();
841 }
842
843 assert_json_eq(
844 &buf,
845 r#"{"nanos":"11-13-2018 +00:00","micros":"11-13-2018 +00:00","millis":"11-13-2018 +00:00","secs":"11-13-2018 +00:00","name":"a"}
846{"name":"b"}
847"#,
848 );
849 }
850
851 #[test]
852 fn write_dates() {
853 let ts_string = "2018-11-13T17:11:10.011375885995";
854 let ts_millis = ts_string
855 .parse::<chrono::NaiveDateTime>()
856 .unwrap()
857 .and_utc()
858 .timestamp_millis();
859
860 let arr_date32 = Date32Array::from(vec![
861 Some(i32::try_from(ts_millis / 1000 / (60 * 60 * 24)).unwrap()),
862 None,
863 ]);
864 let arr_date64 = Date64Array::from(vec![Some(ts_millis), None]);
865 let arr_names = StringArray::from(vec![Some("a"), Some("b")]);
866
867 let schema = Schema::new(vec![
868 Field::new("date32", arr_date32.data_type().clone(), true),
869 Field::new("date64", arr_date64.data_type().clone(), true),
870 Field::new("name", arr_names.data_type().clone(), false),
871 ]);
872 let schema = Arc::new(schema);
873
874 let batch = RecordBatch::try_new(
875 schema,
876 vec![
877 Arc::new(arr_date32),
878 Arc::new(arr_date64),
879 Arc::new(arr_names),
880 ],
881 )
882 .unwrap();
883
884 let mut buf = Vec::new();
885 {
886 let mut writer = LineDelimitedWriter::new(&mut buf);
887 writer.write_batches(&[&batch]).unwrap();
888 }
889
890 assert_json_eq(
891 &buf,
892 r#"{"date32":"2018-11-13","date64":"2018-11-13T17:11:10.011","name":"a"}
893{"name":"b"}
894"#,
895 );
896
897 let mut buf = Vec::new();
898 {
899 let mut writer = WriterBuilder::new()
900 .with_date_format("%m-%d-%Y".to_string())
901 .with_datetime_format("%m-%d-%Y %Mmin %Ssec %Hhour".to_string())
902 .build::<_, LineDelimited>(&mut buf);
903 writer.write_batches(&[&batch]).unwrap();
904 }
905
906 assert_json_eq(
907 &buf,
908 r#"{"date32":"11-13-2018","date64":"11-13-2018 11min 10sec 17hour","name":"a"}
909{"name":"b"}
910"#,
911 );
912 }
913
914 #[test]
915 fn write_times() {
916 let arr_time32sec = Time32SecondArray::from(vec![Some(120), None]);
917 let arr_time32msec = Time32MillisecondArray::from(vec![Some(120), None]);
918 let arr_time64usec = Time64MicrosecondArray::from(vec![Some(120), None]);
919 let arr_time64nsec = Time64NanosecondArray::from(vec![Some(120), None]);
920 let arr_names = StringArray::from(vec![Some("a"), Some("b")]);
921
922 let schema = Schema::new(vec![
923 Field::new("time32sec", arr_time32sec.data_type().clone(), true),
924 Field::new("time32msec", arr_time32msec.data_type().clone(), true),
925 Field::new("time64usec", arr_time64usec.data_type().clone(), true),
926 Field::new("time64nsec", arr_time64nsec.data_type().clone(), true),
927 Field::new("name", arr_names.data_type().clone(), true),
928 ]);
929 let schema = Arc::new(schema);
930
931 let batch = RecordBatch::try_new(
932 schema,
933 vec![
934 Arc::new(arr_time32sec),
935 Arc::new(arr_time32msec),
936 Arc::new(arr_time64usec),
937 Arc::new(arr_time64nsec),
938 Arc::new(arr_names),
939 ],
940 )
941 .unwrap();
942
943 let mut buf = Vec::new();
944 {
945 let mut writer = LineDelimitedWriter::new(&mut buf);
946 writer.write_batches(&[&batch]).unwrap();
947 }
948
949 assert_json_eq(
950 &buf,
951 r#"{"time32sec":"00:02:00","time32msec":"00:00:00.120","time64usec":"00:00:00.000120","time64nsec":"00:00:00.000000120","name":"a"}
952{"name":"b"}
953"#,
954 );
955
956 let mut buf = Vec::new();
957 {
958 let mut writer = WriterBuilder::new()
959 .with_time_format("%H-%M-%S %f".to_string())
960 .build::<_, LineDelimited>(&mut buf);
961 writer.write_batches(&[&batch]).unwrap();
962 }
963
964 assert_json_eq(
965 &buf,
966 r#"{"time32sec":"00-02-00 000000000","time32msec":"00-00-00 120000000","time64usec":"00-00-00 000120000","time64nsec":"00-00-00 000000120","name":"a"}
967{"name":"b"}
968"#,
969 );
970 }
971
972 #[test]
973 fn write_durations() {
974 let arr_durationsec = DurationSecondArray::from(vec![Some(120), None]);
975 let arr_durationmsec = DurationMillisecondArray::from(vec![Some(120), None]);
976 let arr_durationusec = DurationMicrosecondArray::from(vec![Some(120), None]);
977 let arr_durationnsec = DurationNanosecondArray::from(vec![Some(120), None]);
978 let arr_names = StringArray::from(vec![Some("a"), Some("b")]);
979
980 let schema = Schema::new(vec![
981 Field::new("duration_sec", arr_durationsec.data_type().clone(), true),
982 Field::new("duration_msec", arr_durationmsec.data_type().clone(), true),
983 Field::new("duration_usec", arr_durationusec.data_type().clone(), true),
984 Field::new("duration_nsec", arr_durationnsec.data_type().clone(), true),
985 Field::new("name", arr_names.data_type().clone(), true),
986 ]);
987 let schema = Arc::new(schema);
988
989 let batch = RecordBatch::try_new(
990 schema,
991 vec![
992 Arc::new(arr_durationsec),
993 Arc::new(arr_durationmsec),
994 Arc::new(arr_durationusec),
995 Arc::new(arr_durationnsec),
996 Arc::new(arr_names),
997 ],
998 )
999 .unwrap();
1000
1001 let mut buf = Vec::new();
1002 {
1003 let mut writer = LineDelimitedWriter::new(&mut buf);
1004 writer.write_batches(&[&batch]).unwrap();
1005 }
1006
1007 assert_json_eq(
1008 &buf,
1009 r#"{"duration_sec":"PT120S","duration_msec":"PT0.12S","duration_usec":"PT0.00012S","duration_nsec":"PT0.00000012S","name":"a"}
1010{"name":"b"}
1011"#,
1012 );
1013 }
1014
1015 #[test]
1016 fn write_nested_structs() {
1017 let schema = Schema::new(vec![
1018 Field::new(
1019 "c1",
1020 DataType::Struct(Fields::from(vec![
1021 Field::new("c11", DataType::Int32, true),
1022 Field::new(
1023 "c12",
1024 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
1025 false,
1026 ),
1027 ])),
1028 false,
1029 ),
1030 Field::new("c2", DataType::Utf8, false),
1031 ]);
1032
1033 let c1 = StructArray::from(vec![
1034 (
1035 Arc::new(Field::new("c11", DataType::Int32, true)),
1036 Arc::new(Int32Array::from(vec![Some(1), None, Some(5)])) as ArrayRef,
1037 ),
1038 (
1039 Arc::new(Field::new(
1040 "c12",
1041 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
1042 false,
1043 )),
1044 Arc::new(StructArray::from(vec![(
1045 Arc::new(Field::new("c121", DataType::Utf8, false)),
1046 Arc::new(StringArray::from(vec![Some("e"), Some("f"), Some("g")])) as ArrayRef,
1047 )])) as ArrayRef,
1048 ),
1049 ]);
1050 let c2 = StringArray::from(vec![Some("a"), Some("b"), Some("c")]);
1051
1052 let batch =
1053 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap();
1054
1055 let mut buf = Vec::new();
1056 {
1057 let mut writer = LineDelimitedWriter::new(&mut buf);
1058 writer.write_batches(&[&batch]).unwrap();
1059 }
1060
1061 assert_json_eq(
1062 &buf,
1063 r#"{"c1":{"c11":1,"c12":{"c121":"e"}},"c2":"a"}
1064{"c1":{"c12":{"c121":"f"}},"c2":"b"}
1065{"c1":{"c11":5,"c12":{"c121":"g"}},"c2":"c"}
1066"#,
1067 );
1068 }
1069
1070 #[test]
1071 fn write_struct_with_list_field() {
1072 let field_c_list = Arc::new(Field::new("c_list", DataType::Utf8, false));
1073 let field_c1 = Field::new("c1", DataType::List(field_c_list.clone()), false);
1074 let field_c2 = Field::new("c2", DataType::Int32, false);
1075 let schema = Schema::new(vec![field_c1.clone(), field_c2]);
1076
1077 let a_values = StringArray::from(vec!["a", "a1", "b", "c", "d", "e"]);
1078 let a = ListArray::new(
1080 field_c_list,
1081 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 4, 5, 6])),
1082 Arc::new(a_values),
1083 None,
1084 );
1085
1086 let b = Int32Array::from(vec![1, 2, 3, 4, 5]);
1087
1088 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(a), Arc::new(b)]).unwrap();
1089
1090 let mut buf = Vec::new();
1091 {
1092 let mut writer = LineDelimitedWriter::new(&mut buf);
1093 writer.write_batches(&[&batch]).unwrap();
1094 }
1095
1096 assert_json_eq(
1097 &buf,
1098 r#"{"c1":["a","a1"],"c2":1}
1099{"c1":["b"],"c2":2}
1100{"c1":["c"],"c2":3}
1101{"c1":["d"],"c2":4}
1102{"c1":["e"],"c2":5}
1103"#,
1104 );
1105 }
1106
1107 #[test]
1108 fn write_nested_list() {
1109 let field_b = Arc::new(Field::new("b", DataType::Int32, false));
1110 let field_a = Arc::new(Field::new("a", DataType::List(field_b.clone()), false));
1111 let field_c1 = Field::new("c1", DataType::List(field_a.clone()), false);
1112 let field_c2 = Field::new("c2", DataType::Utf8, true);
1113 let schema = Schema::new(vec![field_c1.clone(), field_c2]);
1114
1115 let a_values = Int32Array::from(vec![1, 2, 3, 4, 5, 6]);
1117
1118 let a_list = ListArray::new(
1119 field_b,
1120 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 6])),
1121 Arc::new(a_values),
1122 None,
1123 );
1124
1125 let c1 = ListArray::new(
1126 field_a,
1127 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 2, 3])),
1128 Arc::new(a_list),
1129 None,
1130 );
1131 let c2 = StringArray::from(vec![Some("foo"), Some("bar"), None]);
1132
1133 let batch =
1134 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap();
1135
1136 let mut buf = Vec::new();
1137 {
1138 let mut writer = LineDelimitedWriter::new(&mut buf);
1139 writer.write_batches(&[&batch]).unwrap();
1140 }
1141
1142 assert_json_eq(
1143 &buf,
1144 r#"{"c1":[[1,2],[3]],"c2":"foo"}
1145{"c1":[],"c2":"bar"}
1146{"c1":[[4,5,6]]}
1147"#,
1148 );
1149 }
1150
1151 #[test]
1152 fn write_list_of_struct() {
1153 let field_c1 = Field::new(
1154 "c1",
1155 DataType::List(Arc::new(Field::new(
1156 "s",
1157 DataType::Struct(Fields::from(vec![
1158 Field::new("c11", DataType::Int32, true),
1159 Field::new(
1160 "c12",
1161 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
1162 false,
1163 ),
1164 ])),
1165 false,
1166 ))),
1167 true,
1168 );
1169 let field_c2 = Field::new("c2", DataType::Int32, false);
1170 let schema = Schema::new(vec![field_c1.clone(), field_c2]);
1171
1172 let struct_values = StructArray::from(vec![
1173 (
1174 Arc::new(Field::new("c11", DataType::Int32, true)),
1175 Arc::new(Int32Array::from(vec![Some(1), None, Some(5)])) as ArrayRef,
1176 ),
1177 (
1178 Arc::new(Field::new(
1179 "c12",
1180 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
1181 false,
1182 )),
1183 Arc::new(StructArray::from(vec![(
1184 Arc::new(Field::new("c121", DataType::Utf8, false)),
1185 Arc::new(StringArray::from(vec![Some("e"), Some("f"), Some("g")])) as ArrayRef,
1186 )])) as ArrayRef,
1187 ),
1188 ]);
1189
1190 let c1_inner = match field_c1.data_type() {
1195 DataType::List(f) => f.clone(),
1196 _ => unreachable!(),
1197 };
1198 let c1 = ListArray::new(
1199 c1_inner,
1200 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 2, 3])),
1201 Arc::new(struct_values),
1202 Some(NullBuffer::from(vec![true, false, true])),
1203 );
1204
1205 let c2 = Int32Array::from(vec![1, 2, 3]);
1206
1207 let batch =
1208 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap();
1209
1210 let mut buf = Vec::new();
1211 {
1212 let mut writer = LineDelimitedWriter::new(&mut buf);
1213 writer.write_batches(&[&batch]).unwrap();
1214 }
1215
1216 assert_json_eq(
1217 &buf,
1218 r#"{"c1":[{"c11":1,"c12":{"c121":"e"}},{"c12":{"c121":"f"}}],"c2":1}
1219{"c2":2}
1220{"c1":[{"c11":5,"c12":{"c121":"g"}}],"c2":3}
1221"#,
1222 );
1223 }
1224
1225 fn assert_write_list_view<O: OffsetSizeTrait>() {
1226 let field = Arc::new(Field::new("item", DataType::Int32, true));
1227 let data_type = GenericListViewArray::<O>::DATA_TYPE_CONSTRUCTOR(field.clone());
1228 let schema = Schema::new(vec![Field::new("lv", data_type, true)]);
1229
1230 let values = Int32Array::from(vec![Some(1), Some(2), Some(3), Some(4), None, Some(6)]);
1232 let offsets = [0, 3, 0, 5]
1233 .iter()
1234 .map(|&v| O::from_usize(v).unwrap())
1235 .collect::<Vec<_>>();
1236 let sizes = [3, 2, 0, 1]
1237 .iter()
1238 .map(|&v| O::from_usize(v).unwrap())
1239 .collect::<Vec<_>>();
1240 let list_view = GenericListViewArray::<O>::try_new(
1241 field,
1242 ScalarBuffer::from(offsets),
1243 ScalarBuffer::from(sizes),
1244 Arc::new(values),
1245 Some(NullBuffer::from_iter([true, true, false, true])),
1246 )
1247 .unwrap();
1248
1249 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(list_view)]).unwrap();
1250
1251 let mut buf = Vec::new();
1252 {
1253 let mut writer = LineDelimitedWriter::new(&mut buf);
1254 writer.write_batches(&[&batch]).unwrap();
1255 }
1256
1257 assert_json_eq(
1258 &buf,
1259 r#"{"lv":[1,2,3]}
1260{"lv":[4,null]}
1261{}
1262{"lv":[6]}
1263"#,
1264 );
1265 }
1266
1267 #[test]
1268 fn write_list_view() {
1269 assert_write_list_view::<i32>();
1270 assert_write_list_view::<i64>();
1271 }
1272
1273 fn test_write_for_file(test_file: &str, remove_nulls: bool) {
1274 let file = File::open(test_file).unwrap();
1275 let mut reader = BufReader::new(file);
1276 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
1277 reader.rewind().unwrap();
1278
1279 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(1024);
1280 let mut reader = builder.build(reader).unwrap();
1281 let batch = reader.next().unwrap().unwrap();
1282
1283 let mut buf = Vec::new();
1284 {
1285 if remove_nulls {
1286 let mut writer = LineDelimitedWriter::new(&mut buf);
1287 writer.write_batches(&[&batch]).unwrap();
1288 } else {
1289 let mut writer = WriterBuilder::new()
1290 .with_explicit_nulls(true)
1291 .build::<_, LineDelimited>(&mut buf);
1292 writer.write_batches(&[&batch]).unwrap();
1293 }
1294 }
1295
1296 let result = str::from_utf8(&buf).unwrap();
1297 let expected = read_to_string(test_file).unwrap();
1298 for (r, e) in result.lines().zip(expected.lines()) {
1299 let mut expected_json = serde_json::from_str::<Value>(e).unwrap();
1300 if remove_nulls {
1301 if let Value::Object(obj) = expected_json {
1303 expected_json =
1304 Value::Object(obj.into_iter().filter(|(_, v)| *v != Value::Null).collect());
1305 }
1306 }
1307 assert_eq!(serde_json::from_str::<Value>(r).unwrap(), expected_json,);
1308 }
1309 }
1310
1311 #[test]
1312 fn write_basic_rows() {
1313 test_write_for_file("test/data/basic.json", true);
1314 }
1315
1316 #[test]
1317 fn write_arrays() {
1318 test_write_for_file("test/data/arrays.json", true);
1319 }
1320
1321 #[test]
1322 fn write_basic_nulls() {
1323 test_write_for_file("test/data/basic_nulls.json", true);
1324 }
1325
1326 #[test]
1327 fn write_nested_with_nulls() {
1328 test_write_for_file("test/data/nested_with_nulls.json", false);
1329 }
1330
1331 #[test]
1332 fn json_line_writer_empty() {
1333 let mut writer = LineDelimitedWriter::new(vec![] as Vec<u8>);
1334 writer.finish().unwrap();
1335 assert_eq!(str::from_utf8(&writer.into_inner()).unwrap(), "");
1336 }
1337
1338 #[test]
1339 fn json_array_writer_empty() {
1340 let mut writer = ArrayWriter::new(vec![] as Vec<u8>);
1341 writer.finish().unwrap();
1342 assert_eq!(str::from_utf8(&writer.into_inner()).unwrap(), "[]");
1343 }
1344
1345 #[test]
1346 fn json_line_writer_empty_batch() {
1347 let mut writer = LineDelimitedWriter::new(vec![] as Vec<u8>);
1348
1349 let array = Int32Array::from(Vec::<i32>::new());
1350 let schema = Schema::new(vec![Field::new("c", DataType::Int32, true)]);
1351 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
1352
1353 writer.write(&batch).unwrap();
1354 writer.finish().unwrap();
1355 assert_eq!(str::from_utf8(&writer.into_inner()).unwrap(), "");
1356 }
1357
1358 #[test]
1359 fn json_array_writer_empty_batch() {
1360 let mut writer = ArrayWriter::new(vec![] as Vec<u8>);
1361
1362 let array = Int32Array::from(Vec::<i32>::new());
1363 let schema = Schema::new(vec![Field::new("c", DataType::Int32, true)]);
1364 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
1365
1366 writer.write(&batch).unwrap();
1367 writer.finish().unwrap();
1368 assert_eq!(str::from_utf8(&writer.into_inner()).unwrap(), "[]");
1369 }
1370
1371 #[test]
1372 fn json_struct_array_nulls() {
1373 let inner = ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
1374 Some(vec![Some(1), Some(2)]),
1375 Some(vec![None]),
1376 Some(vec![]),
1377 Some(vec![Some(3), None]), Some(vec![Some(4), Some(5)]),
1379 None, None,
1381 ]);
1382
1383 let field = Arc::new(Field::new("list", inner.data_type().clone(), true));
1384 let array = Arc::new(inner) as ArrayRef;
1385 let struct_array_a = StructArray::from((
1386 vec![(field.clone(), array.clone())],
1387 Buffer::from([0b01010111]),
1388 ));
1389 let struct_array_b = StructArray::from(vec![(field, array)]);
1390
1391 let schema = Schema::new(vec![
1392 Field::new_struct("a", struct_array_a.fields().clone(), true),
1393 Field::new_struct("b", struct_array_b.fields().clone(), true),
1394 ]);
1395
1396 let batch = RecordBatch::try_new(
1397 Arc::new(schema),
1398 vec![Arc::new(struct_array_a), Arc::new(struct_array_b)],
1399 )
1400 .unwrap();
1401
1402 let mut buf = Vec::new();
1403 {
1404 let mut writer = LineDelimitedWriter::new(&mut buf);
1405 writer.write_batches(&[&batch]).unwrap();
1406 }
1407
1408 assert_json_eq(
1409 &buf,
1410 r#"{"a":{"list":[1,2]},"b":{"list":[1,2]}}
1411{"a":{"list":[null]},"b":{"list":[null]}}
1412{"a":{"list":[]},"b":{"list":[]}}
1413{"b":{"list":[3,null]}}
1414{"a":{"list":[4,5]},"b":{"list":[4,5]}}
1415{"b":{}}
1416{"a":{},"b":{}}
1417"#,
1418 );
1419 }
1420
1421 fn run_json_writer_map_with_keys(keys_array: ArrayRef) {
1422 let values_array = super::Int64Array::from(vec![10, 20, 30, 40, 50]);
1423
1424 let keys_field = Arc::new(Field::new(
1425 Field::MAP_KEY_FIELD_DEFAULT_NAME,
1426 keys_array.data_type().clone(),
1427 false,
1428 ));
1429 let values_field = Arc::new(Field::new(
1430 Field::MAP_VALUE_FIELD_DEFAULT_NAME,
1431 DataType::Int64,
1432 false,
1433 ));
1434 let entry_struct = StructArray::from(vec![
1435 (keys_field, keys_array.clone()),
1436 (values_field, Arc::new(values_array) as ArrayRef),
1437 ]);
1438
1439 let entries_field = Arc::new(Field::new(
1440 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1441 entry_struct.data_type().clone(),
1442 false,
1443 ));
1444
1445 let map = MapArray::new(
1447 entries_field.clone(),
1448 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 1, 1, 1, 4, 5, 5])),
1449 entry_struct,
1450 Some(NullBuffer::from(vec![true, false, true, true, true, true])),
1451 false,
1452 );
1453
1454 let map_field = Field::new("map", DataType::Map(entries_field, false), true);
1455 let schema = Arc::new(Schema::new(vec![map_field]));
1456
1457 let batch = RecordBatch::try_new(schema, vec![Arc::new(map)]).unwrap();
1458
1459 let mut buf = Vec::new();
1460 {
1461 let mut writer = LineDelimitedWriter::new(&mut buf);
1462 writer.write_batches(&[&batch]).unwrap();
1463 }
1464
1465 assert_json_eq(
1466 &buf,
1467 r#"{"map":{"foo":10}}
1468{}
1469{"map":{}}
1470{"map":{"bar":20,"baz":30,"qux":40}}
1471{"map":{"quux":50}}
1472{"map":{}}
1473"#,
1474 );
1475 }
1476
1477 #[test]
1478 fn json_writer_map() {
1479 let keys_utf8 = super::StringArray::from(vec!["foo", "bar", "baz", "qux", "quux"]);
1481 run_json_writer_map_with_keys(Arc::new(keys_utf8) as ArrayRef);
1482
1483 let keys_large = super::LargeStringArray::from(vec!["foo", "bar", "baz", "qux", "quux"]);
1485 run_json_writer_map_with_keys(Arc::new(keys_large) as ArrayRef);
1486
1487 let keys_view = super::StringViewArray::from(vec!["foo", "bar", "baz", "qux", "quux"]);
1489 run_json_writer_map_with_keys(Arc::new(keys_view) as ArrayRef);
1490 }
1491
1492 #[test]
1493 fn test_write_single_batch() {
1494 let test_file = "test/data/basic.json";
1495 let file = File::open(test_file).unwrap();
1496 let mut reader = BufReader::new(file);
1497 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
1498 reader.rewind().unwrap();
1499
1500 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(1024);
1501 let mut reader = builder.build(reader).unwrap();
1502 let batch = reader.next().unwrap().unwrap();
1503
1504 let mut buf = Vec::new();
1505 {
1506 let mut writer = LineDelimitedWriter::new(&mut buf);
1507 writer.write(&batch).unwrap();
1508 }
1509
1510 let result = str::from_utf8(&buf).unwrap();
1511 let expected = read_to_string(test_file).unwrap();
1512 for (r, e) in result.lines().zip(expected.lines()) {
1513 let mut expected_json = serde_json::from_str::<Value>(e).unwrap();
1514 if let Value::Object(obj) = expected_json {
1516 expected_json =
1517 Value::Object(obj.into_iter().filter(|(_, v)| *v != Value::Null).collect());
1518 }
1519 assert_eq!(serde_json::from_str::<Value>(r).unwrap(), expected_json,);
1520 }
1521 }
1522
1523 #[test]
1524 fn test_write_multi_batches() {
1525 let test_file = "test/data/basic.json";
1526
1527 let schema = SchemaRef::new(Schema::new(vec![
1528 Field::new("a", DataType::Int64, true),
1529 Field::new("b", DataType::Float64, true),
1530 Field::new("c", DataType::Boolean, true),
1531 Field::new("d", DataType::Utf8, true),
1532 Field::new("e", DataType::Utf8, true),
1533 Field::new("f", DataType::Utf8, true),
1534 Field::new("g", DataType::Timestamp(TimeUnit::Millisecond, None), true),
1535 Field::new("h", DataType::Float16, true),
1536 ]));
1537
1538 let mut reader = ReaderBuilder::new(schema.clone())
1539 .build(BufReader::new(File::open(test_file).unwrap()))
1540 .unwrap();
1541 let batch = reader.next().unwrap().unwrap();
1542
1543 let batches = [&RecordBatch::new_empty(schema), &batch, &batch];
1545
1546 let mut buf = Vec::new();
1547 {
1548 let mut writer = LineDelimitedWriter::new(&mut buf);
1549 writer.write_batches(&batches).unwrap();
1550 }
1551
1552 let result = str::from_utf8(&buf).unwrap();
1553 let expected = read_to_string(test_file).unwrap();
1554 let expected = format!("{expected}\n{expected}");
1556 for (r, e) in result.lines().zip(expected.lines()) {
1557 let mut expected_json = serde_json::from_str::<Value>(e).unwrap();
1558 if let Value::Object(obj) = expected_json {
1560 expected_json =
1561 Value::Object(obj.into_iter().filter(|(_, v)| *v != Value::Null).collect());
1562 }
1563 assert_eq!(serde_json::from_str::<Value>(r).unwrap(), expected_json,);
1564 }
1565 }
1566
1567 #[test]
1568 fn test_writer_explicit_nulls() -> Result<(), ArrowError> {
1569 fn nested_list() -> (Arc<ListArray>, Arc<Field>) {
1570 let array = Arc::new(ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
1571 Some(vec![None, None, None]),
1572 Some(vec![Some(1), Some(2), Some(3)]),
1573 None,
1574 Some(vec![None, None, None]),
1575 ]));
1576 let field = Arc::new(Field::new("list", array.data_type().clone(), true));
1577 (array, field)
1579 }
1580
1581 fn nested_dict() -> (Arc<DictionaryArray<Int32Type>>, Arc<Field>) {
1582 let array = Arc::new(DictionaryArray::from_iter(vec![
1583 Some("cupcakes"),
1584 None,
1585 Some("bear"),
1586 Some("kuma"),
1587 ]));
1588 let field = Arc::new(Field::new("dict", array.data_type().clone(), true));
1589 (array, field)
1591 }
1592
1593 fn nested_map() -> (Arc<MapArray>, Arc<Field>) {
1594 let string_builder = StringBuilder::new();
1595 let int_builder = Int64Builder::new();
1596 let mut builder = MapBuilder::new(None, string_builder, int_builder);
1597
1598 builder.keys().append_value("foo");
1600 builder.values().append_value(10);
1601 builder.append(true).unwrap();
1602
1603 builder.append(false).unwrap();
1604
1605 builder.append(true).unwrap();
1606
1607 builder.keys().append_value("bar");
1608 builder.values().append_value(20);
1609 builder.keys().append_value("baz");
1610 builder.values().append_value(30);
1611 builder.keys().append_value("qux");
1612 builder.values().append_value(40);
1613 builder.append(true).unwrap();
1614
1615 let array = Arc::new(builder.finish());
1616 let field = Arc::new(Field::new("map", array.data_type().clone(), true));
1617 (array, field)
1618 }
1619
1620 fn root_list() -> (Arc<ListArray>, Field) {
1621 let struct_array = StructArray::from(vec![
1622 (
1623 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1624 Arc::new(StringArray::from(vec![Some("a"), Some("b"), None, None])) as ArrayRef,
1625 ),
1626 (
1627 Arc::new(Field::new("int32", DataType::Int32, true)),
1628 Arc::new(Int32Array::from(vec![Some(1), None, Some(5), None])) as ArrayRef,
1629 ),
1630 ]);
1631
1632 let values_field =
1633 Arc::new(Field::new("struct", struct_array.data_type().clone(), true));
1634 let field = Field::new_list("list", values_field.as_ref().clone(), true);
1635
1636 let array = Arc::new(ListArray::new(
1638 values_field,
1639 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 2, 3, 3])),
1640 Arc::new(struct_array),
1641 Some(NullBuffer::from(vec![true, false, true, false])),
1642 ));
1643 (array, field)
1644 }
1645
1646 let (nested_list_array, nested_list_field) = nested_list();
1647 let (nested_dict_array, nested_dict_field) = nested_dict();
1648 let (nested_map_array, nested_map_field) = nested_map();
1649 let (root_list_array, root_list_field) = root_list();
1650
1651 let schema = Schema::new(vec![
1652 Field::new("date", DataType::Date32, true),
1653 Field::new("null", DataType::Null, true),
1654 Field::new_struct(
1655 "struct",
1656 vec![
1657 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1658 nested_list_field.clone(),
1659 nested_dict_field.clone(),
1660 nested_map_field.clone(),
1661 ],
1662 true,
1663 ),
1664 root_list_field,
1665 ]);
1666
1667 let arr_date32 = Date32Array::from(vec![Some(0), None, Some(1), None]);
1668 let arr_null = NullArray::new(4);
1669 let arr_struct = StructArray::from(vec![
1670 (
1672 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1673 Arc::new(StringArray::from(vec![Some("a"), None, None, Some("b")])) as ArrayRef,
1674 ),
1675 (nested_list_field, nested_list_array as ArrayRef),
1677 (nested_dict_field, nested_dict_array as ArrayRef),
1679 (nested_map_field, nested_map_array as ArrayRef),
1681 ]);
1682
1683 let batch = RecordBatch::try_new(
1684 Arc::new(schema),
1685 vec![
1686 Arc::new(arr_date32),
1688 Arc::new(arr_null),
1690 Arc::new(arr_struct),
1691 root_list_array,
1693 ],
1694 )?;
1695
1696 let mut buf = Vec::new();
1697 {
1698 let mut writer = WriterBuilder::new()
1699 .with_explicit_nulls(true)
1700 .build::<_, JsonArray>(&mut buf);
1701 writer.write_batches(&[&batch])?;
1702 writer.finish()?;
1703 }
1704
1705 let actual = serde_json::from_slice::<Vec<Value>>(&buf).unwrap();
1706 let expected = serde_json::from_value::<Vec<Value>>(json!([
1707 {
1708 "date": "1970-01-01",
1709 "list": [
1710 {
1711 "int32": 1,
1712 "utf8": "a"
1713 },
1714 {
1715 "int32": null,
1716 "utf8": "b"
1717 }
1718 ],
1719 "null": null,
1720 "struct": {
1721 "dict": "cupcakes",
1722 "list": [
1723 null,
1724 null,
1725 null
1726 ],
1727 "map": {
1728 "foo": 10
1729 },
1730 "utf8": "a"
1731 }
1732 },
1733 {
1734 "date": null,
1735 "list": null,
1736 "null": null,
1737 "struct": {
1738 "dict": null,
1739 "list": [
1740 1,
1741 2,
1742 3
1743 ],
1744 "map": null,
1745 "utf8": null
1746 }
1747 },
1748 {
1749 "date": "1970-01-02",
1750 "list": [
1751 {
1752 "int32": 5,
1753 "utf8": null
1754 }
1755 ],
1756 "null": null,
1757 "struct": {
1758 "dict": "bear",
1759 "list": null,
1760 "map": {},
1761 "utf8": null
1762 }
1763 },
1764 {
1765 "date": null,
1766 "list": null,
1767 "null": null,
1768 "struct": {
1769 "dict": "kuma",
1770 "list": [
1771 null,
1772 null,
1773 null
1774 ],
1775 "map": {
1776 "bar": 20,
1777 "baz": 30,
1778 "qux": 40
1779 },
1780 "utf8": "b"
1781 }
1782 }
1783 ]))
1784 .unwrap();
1785
1786 assert_eq!(actual, expected);
1787
1788 Ok(())
1789 }
1790
1791 fn build_array_binary<O: OffsetSizeTrait>(values: &[Option<&[u8]>]) -> RecordBatch {
1792 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1793 "bytes",
1794 GenericBinaryType::<O>::DATA_TYPE,
1795 true,
1796 )]));
1797 let mut builder = GenericByteBuilder::<GenericBinaryType<O>>::new();
1798 for value in values {
1799 match value {
1800 Some(v) => builder.append_value(v),
1801 None => builder.append_null(),
1802 }
1803 }
1804 let array = Arc::new(builder.finish()) as ArrayRef;
1805 RecordBatch::try_new(schema, vec![array]).unwrap()
1806 }
1807
1808 fn build_array_binary_view(values: &[Option<&[u8]>]) -> RecordBatch {
1809 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1810 "bytes",
1811 DataType::BinaryView,
1812 true,
1813 )]));
1814 let mut builder = BinaryViewBuilder::new();
1815 for value in values {
1816 match value {
1817 Some(v) => builder.append_value(v),
1818 None => builder.append_null(),
1819 }
1820 }
1821 let array = Arc::new(builder.finish()) as ArrayRef;
1822 RecordBatch::try_new(schema, vec![array]).unwrap()
1823 }
1824
1825 fn assert_binary_json(batch: &RecordBatch) {
1826 {
1828 let mut buf = Vec::new();
1829 let json_value: Value = {
1830 let mut writer = WriterBuilder::new()
1831 .with_explicit_nulls(true)
1832 .build::<_, JsonArray>(&mut buf);
1833 writer.write(batch).unwrap();
1834 writer.close().unwrap();
1835 serde_json::from_slice(&buf).unwrap()
1836 };
1837
1838 assert_eq!(
1839 json!([
1840 {
1841 "bytes": "4e656420466c616e64657273"
1842 },
1843 {
1844 "bytes": null },
1846 {
1847 "bytes": "54726f79204d63436c757265"
1848 }
1849 ]),
1850 json_value,
1851 );
1852 }
1853
1854 {
1856 let mut buf = Vec::new();
1857 let json_value: Value = {
1858 let mut writer = ArrayWriter::new(&mut buf);
1861 writer.write(batch).unwrap();
1862 writer.close().unwrap();
1863 serde_json::from_slice(&buf).unwrap()
1864 };
1865
1866 assert_eq!(
1867 json!([
1868 { "bytes": "4e656420466c616e64657273" },
1869 {},
1870 { "bytes": "54726f79204d63436c757265" }
1871 ]),
1872 json_value
1873 );
1874 }
1875 }
1876
1877 #[test]
1878 fn test_writer_binary() {
1879 let values: [Option<&[u8]>; 3] = [
1880 Some(b"Ned Flanders" as &[u8]),
1881 None,
1882 Some(b"Troy McClure" as &[u8]),
1883 ];
1884 {
1886 let batch = build_array_binary::<i32>(&values);
1887 assert_binary_json(&batch);
1888 }
1889 {
1891 let batch = build_array_binary::<i64>(&values);
1892 assert_binary_json(&batch);
1893 }
1894 {
1895 let batch = build_array_binary_view(&values);
1896 assert_binary_json(&batch);
1897 }
1898 }
1899
1900 #[test]
1901 fn test_writer_fixed_size_binary() {
1902 let size = 11;
1904 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1905 "bytes",
1906 DataType::FixedSizeBinary(size),
1907 true,
1908 )]));
1909
1910 let mut builder = FixedSizeBinaryBuilder::new(size);
1912 let values = [Some(b"hello world"), None, Some(b"summer rain")];
1913 for value in values {
1914 match value {
1915 Some(v) => builder.append_value(v).unwrap(),
1916 None => builder.append_null(),
1917 }
1918 }
1919 let array = Arc::new(builder.finish()) as ArrayRef;
1920 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
1921
1922 {
1924 let mut buf = Vec::new();
1925 let json_value: Value = {
1926 let mut writer = WriterBuilder::new()
1927 .with_explicit_nulls(true)
1928 .build::<_, JsonArray>(&mut buf);
1929 writer.write(&batch).unwrap();
1930 writer.close().unwrap();
1931 serde_json::from_slice(&buf).unwrap()
1932 };
1933
1934 assert_eq!(
1935 json!([
1936 {
1937 "bytes": "68656c6c6f20776f726c64"
1938 },
1939 {
1940 "bytes": null },
1942 {
1943 "bytes": "73756d6d6572207261696e"
1944 }
1945 ]),
1946 json_value,
1947 );
1948 }
1949 {
1951 let mut buf = Vec::new();
1952 let json_value: Value = {
1953 let mut writer = ArrayWriter::new(&mut buf);
1956 writer.write(&batch).unwrap();
1957 writer.close().unwrap();
1958 serde_json::from_slice(&buf).unwrap()
1959 };
1960
1961 assert_eq!(
1962 json!([
1963 {
1964 "bytes": "68656c6c6f20776f726c64"
1965 },
1966 {}, {
1968 "bytes": "73756d6d6572207261696e"
1969 }
1970 ]),
1971 json_value,
1972 );
1973 }
1974 }
1975
1976 #[test]
1977 fn test_writer_fixed_size_list() {
1978 let size = 3;
1979 let field = FieldRef::new(Field::new_list_field(DataType::Int32, true));
1980 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1981 "list",
1982 DataType::FixedSizeList(field, size),
1983 true,
1984 )]));
1985
1986 let values_builder = Int32Builder::new();
1987 let mut list_builder = FixedSizeListBuilder::new(values_builder, size);
1988 let lists = [
1989 Some([Some(1), Some(2), None]),
1990 Some([Some(3), None, Some(4)]),
1991 Some([None, Some(5), Some(6)]),
1992 None,
1993 ];
1994 for list in lists {
1995 match list {
1996 Some(l) => {
1997 for value in l {
1998 match value {
1999 Some(v) => list_builder.values().append_value(v),
2000 None => list_builder.values().append_null(),
2001 }
2002 }
2003 list_builder.append(true);
2004 }
2005 None => {
2006 for _ in 0..size {
2007 list_builder.values().append_null();
2008 }
2009 list_builder.append(false);
2010 }
2011 }
2012 }
2013 let array = Arc::new(list_builder.finish()) as ArrayRef;
2014 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
2015
2016 {
2018 let json_value: Value = {
2019 let mut buf = Vec::new();
2020 let mut writer = WriterBuilder::new()
2021 .with_explicit_nulls(true)
2022 .build::<_, JsonArray>(&mut buf);
2023 writer.write(&batch).unwrap();
2024 writer.close().unwrap();
2025 serde_json::from_slice(&buf).unwrap()
2026 };
2027 assert_eq!(
2028 json!([
2029 {"list": [1, 2, null]},
2030 {"list": [3, null, 4]},
2031 {"list": [null, 5, 6]},
2032 {"list": null},
2033 ]),
2034 json_value
2035 );
2036 }
2037 {
2039 let json_value: Value = {
2040 let mut buf = Vec::new();
2041 let mut writer = ArrayWriter::new(&mut buf);
2042 writer.write(&batch).unwrap();
2043 writer.close().unwrap();
2044 serde_json::from_slice(&buf).unwrap()
2045 };
2046 assert_eq!(
2047 json!([
2048 {"list": [1, 2, null]},
2049 {"list": [3, null, 4]},
2050 {"list": [null, 5, 6]},
2051 {}, ]),
2053 json_value
2054 );
2055 }
2056 }
2057
2058 #[test]
2059 fn test_writer_null_dict() {
2060 let keys = Int32Array::from_iter(vec![Some(0), None, Some(1)]);
2061 let values = Arc::new(StringArray::from_iter(vec![Some("a"), None]));
2062 let dict = DictionaryArray::new(keys, values);
2063
2064 let schema = SchemaRef::new(Schema::new(vec![Field::new(
2065 "my_dict",
2066 DataType::Dictionary(DataType::Int32.into(), DataType::Utf8.into()),
2067 true,
2068 )]));
2069
2070 let array = Arc::new(dict) as ArrayRef;
2071 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
2072
2073 let mut json = Vec::new();
2074 let write_builder = WriterBuilder::new().with_explicit_nulls(true);
2075 let mut writer = write_builder.build::<_, JsonArray>(&mut json);
2076 writer.write(&batch).unwrap();
2077 writer.close().unwrap();
2078
2079 let json_str = str::from_utf8(&json).unwrap();
2080 assert_eq!(
2081 json_str,
2082 r#"[{"my_dict":"a"},{"my_dict":null},{"my_dict":""}]"#
2083 )
2084 }
2085
2086 #[test]
2087 fn test_decimal32_encoder() {
2088 let array = Decimal32Array::from_iter_values([1234, 5678, 9012])
2089 .with_precision_and_scale(8, 2)
2090 .unwrap();
2091 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2092 let schema = Schema::new(vec![field]);
2093 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2094
2095 let mut buf = Vec::new();
2096 {
2097 let mut writer = LineDelimitedWriter::new(&mut buf);
2098 writer.write_batches(&[&batch]).unwrap();
2099 }
2100
2101 assert_json_eq(
2102 &buf,
2103 r#"{"decimal":12.34}
2104{"decimal":56.78}
2105{"decimal":90.12}
2106"#,
2107 );
2108 }
2109
2110 #[test]
2111 fn test_decimal64_encoder() {
2112 let array = Decimal64Array::from_iter_values([1234, 5678, 9012])
2113 .with_precision_and_scale(10, 2)
2114 .unwrap();
2115 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2116 let schema = Schema::new(vec![field]);
2117 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2118
2119 let mut buf = Vec::new();
2120 {
2121 let mut writer = LineDelimitedWriter::new(&mut buf);
2122 writer.write_batches(&[&batch]).unwrap();
2123 }
2124
2125 assert_json_eq(
2126 &buf,
2127 r#"{"decimal":12.34}
2128{"decimal":56.78}
2129{"decimal":90.12}
2130"#,
2131 );
2132 }
2133
2134 #[test]
2135 fn test_decimal128_encoder() {
2136 let array = Decimal128Array::from_iter_values([1234, 5678, 9012])
2137 .with_precision_and_scale(10, 2)
2138 .unwrap();
2139 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2140 let schema = Schema::new(vec![field]);
2141 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2142
2143 let mut buf = Vec::new();
2144 {
2145 let mut writer = LineDelimitedWriter::new(&mut buf);
2146 writer.write_batches(&[&batch]).unwrap();
2147 }
2148
2149 assert_json_eq(
2150 &buf,
2151 r#"{"decimal":12.34}
2152{"decimal":56.78}
2153{"decimal":90.12}
2154"#,
2155 );
2156 }
2157
2158 #[test]
2159 fn test_decimal256_encoder() {
2160 let array = Decimal256Array::from_iter_values([
2161 i256::from(123400),
2162 i256::from(567800),
2163 i256::from(901200),
2164 ])
2165 .with_precision_and_scale(10, 4)
2166 .unwrap();
2167 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2168 let schema = Schema::new(vec![field]);
2169 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2170
2171 let mut buf = Vec::new();
2172 {
2173 let mut writer = LineDelimitedWriter::new(&mut buf);
2174 writer.write_batches(&[&batch]).unwrap();
2175 }
2176
2177 assert_json_eq(
2178 &buf,
2179 r#"{"decimal":12.3400}
2180{"decimal":56.7800}
2181{"decimal":90.1200}
2182"#,
2183 );
2184 }
2185
2186 #[test]
2187 fn test_decimal_encoder_with_nulls() {
2188 let array = Decimal128Array::from_iter([Some(1234), None, Some(5678)])
2189 .with_precision_and_scale(10, 2)
2190 .unwrap();
2191 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2192 let schema = Schema::new(vec![field]);
2193 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2194
2195 let mut buf = Vec::new();
2196 {
2197 let mut writer = LineDelimitedWriter::new(&mut buf);
2198 writer.write_batches(&[&batch]).unwrap();
2199 }
2200
2201 assert_json_eq(
2202 &buf,
2203 r#"{"decimal":12.34}
2204{}
2205{"decimal":56.78}
2206"#,
2207 );
2208 }
2209
2210 #[test]
2211 fn write_structs_as_list() {
2212 let schema = Schema::new(vec![
2213 Field::new(
2214 "c1",
2215 DataType::Struct(Fields::from(vec![
2216 Field::new("c11", DataType::Int32, true),
2217 Field::new(
2218 "c12",
2219 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
2220 false,
2221 ),
2222 ])),
2223 false,
2224 ),
2225 Field::new("c2", DataType::Utf8, false),
2226 ]);
2227
2228 let c1 = StructArray::from(vec![
2229 (
2230 Arc::new(Field::new("c11", DataType::Int32, true)),
2231 Arc::new(Int32Array::from(vec![Some(1), None, Some(5)])) as ArrayRef,
2232 ),
2233 (
2234 Arc::new(Field::new(
2235 "c12",
2236 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
2237 false,
2238 )),
2239 Arc::new(StructArray::from(vec![(
2240 Arc::new(Field::new("c121", DataType::Utf8, false)),
2241 Arc::new(StringArray::from(vec![Some("e"), Some("f"), Some("g")])) as ArrayRef,
2242 )])) as ArrayRef,
2243 ),
2244 ]);
2245 let c2 = StringArray::from(vec![Some("a"), Some("b"), Some("c")]);
2246
2247 let batch =
2248 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap();
2249
2250 let expected = r#"[[1,["e"]],"a"]
2251[[null,["f"]],"b"]
2252[[5,["g"]],"c"]
2253"#;
2254
2255 let mut buf = Vec::new();
2256 {
2257 let builder = WriterBuilder::new()
2258 .with_explicit_nulls(true)
2259 .with_struct_mode(StructMode::ListOnly);
2260 let mut writer = builder.build::<_, LineDelimited>(&mut buf);
2261 writer.write_batches(&[&batch]).unwrap();
2262 }
2263 assert_json_eq(&buf, expected);
2264
2265 let mut buf = Vec::new();
2266 {
2267 let builder = WriterBuilder::new()
2268 .with_explicit_nulls(false)
2269 .with_struct_mode(StructMode::ListOnly);
2270 let mut writer = builder.build::<_, LineDelimited>(&mut buf);
2271 writer.write_batches(&[&batch]).unwrap();
2272 }
2273 assert_json_eq(&buf, expected);
2274 }
2275
2276 fn make_fallback_encoder_test_data() -> (RecordBatch, Arc<dyn EncoderFactory>) {
2277 #[derive(Debug)]
2280 enum UnionValue {
2281 Int32(i32),
2282 String(String),
2283 }
2284
2285 #[derive(Debug)]
2286 struct UnionEncoder {
2287 array: Vec<Option<UnionValue>>,
2288 }
2289
2290 impl Encoder for UnionEncoder {
2291 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2292 match &self.array[idx] {
2293 None => out.extend_from_slice(b"null"),
2294 Some(UnionValue::Int32(v)) => out.extend_from_slice(v.to_string().as_bytes()),
2295 Some(UnionValue::String(v)) => {
2296 out.extend_from_slice(format!("\"{v}\"").as_bytes())
2297 }
2298 }
2299 }
2300 }
2301
2302 #[derive(Debug)]
2303 struct UnionEncoderFactory;
2304
2305 impl EncoderFactory for UnionEncoderFactory {
2306 fn make_default_encoder<'a>(
2307 &self,
2308 _field: &'a FieldRef,
2309 array: &'a dyn Array,
2310 _options: &'a EncoderOptions,
2311 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2312 let data_type = array.data_type();
2313 let fields = match data_type {
2314 DataType::Union(fields, UnionMode::Sparse) => fields,
2315 _ => return Ok(None),
2316 };
2317 let fields = fields.iter().map(|(_, f)| f).collect::<Vec<_>>();
2319 for f in fields.iter() {
2320 match f.data_type() {
2321 DataType::Null => {}
2322 DataType::Int32 => {}
2323 DataType::Utf8 => {}
2324 _ => return Ok(None),
2325 }
2326 }
2327 let (_, type_ids, _, buffers) = array.as_union().clone().into_parts();
2328 let mut values = Vec::with_capacity(type_ids.len());
2329 for idx in 0..type_ids.len() {
2330 let type_id = type_ids[idx];
2331 let field = &fields[type_id as usize];
2332 let value = match field.data_type() {
2333 DataType::Null => None,
2334 DataType::Int32 => Some(UnionValue::Int32(
2335 buffers[type_id as usize]
2336 .as_primitive::<Int32Type>()
2337 .value(idx),
2338 )),
2339 DataType::Utf8 => Some(UnionValue::String(
2340 buffers[type_id as usize]
2341 .as_string::<i32>()
2342 .value(idx)
2343 .to_string(),
2344 )),
2345 _ => unreachable!(),
2346 };
2347 values.push(value);
2348 }
2349 let array_encoder =
2350 Box::new(UnionEncoder { array: values }) as Box<dyn Encoder + 'a>;
2351 let nulls = array.nulls().cloned();
2352 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2353 }
2354 }
2355
2356 let int_array = Int32Array::from(vec![Some(1), None, None]);
2357 let string_array = StringArray::from(vec![None, Some("a"), None]);
2358 let null_array = NullArray::new(3);
2359 let type_ids = [0_i8, 1, 2].into_iter().collect::<ScalarBuffer<i8>>();
2360
2361 let union_fields = [
2362 (0, Arc::new(Field::new("A", DataType::Int32, false))),
2363 (1, Arc::new(Field::new("B", DataType::Utf8, false))),
2364 (2, Arc::new(Field::new("C", DataType::Null, false))),
2365 ]
2366 .into_iter()
2367 .collect::<UnionFields>();
2368
2369 let children = vec![
2370 Arc::new(int_array) as Arc<dyn Array>,
2371 Arc::new(string_array),
2372 Arc::new(null_array),
2373 ];
2374
2375 let array = UnionArray::try_new(union_fields.clone(), type_ids, None, children).unwrap();
2376
2377 let float_array = Float64Array::from(vec![Some(1.0), None, Some(3.4)]);
2378
2379 let fields = vec![
2380 Field::new(
2381 "union",
2382 DataType::Union(union_fields, UnionMode::Sparse),
2383 true,
2384 ),
2385 Field::new("float", DataType::Float64, true),
2386 ];
2387
2388 let batch = RecordBatch::try_new(
2389 Arc::new(Schema::new(fields)),
2390 vec![
2391 Arc::new(array) as Arc<dyn Array>,
2392 Arc::new(float_array) as Arc<dyn Array>,
2393 ],
2394 )
2395 .unwrap();
2396
2397 (batch, Arc::new(UnionEncoderFactory))
2398 }
2399
2400 #[test]
2401 fn test_fallback_encoder_factory_line_delimited_implicit_nulls() {
2402 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2403
2404 let mut buf = Vec::new();
2405 {
2406 let mut writer = WriterBuilder::new()
2407 .with_encoder_factory(encoder_factory)
2408 .with_explicit_nulls(false)
2409 .build::<_, LineDelimited>(&mut buf);
2410 writer.write_batches(&[&batch]).unwrap();
2411 writer.finish().unwrap();
2412 }
2413
2414 println!("{}", str::from_utf8(&buf).unwrap());
2415
2416 assert_json_eq(
2417 &buf,
2418 r#"{"union":1,"float":1.0}
2419{"union":"a"}
2420{"union":null,"float":3.4}
2421"#,
2422 );
2423 }
2424
2425 #[test]
2426 fn test_fallback_encoder_factory_line_delimited_explicit_nulls() {
2427 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2428
2429 let mut buf = Vec::new();
2430 {
2431 let mut writer = WriterBuilder::new()
2432 .with_encoder_factory(encoder_factory)
2433 .with_explicit_nulls(true)
2434 .build::<_, LineDelimited>(&mut buf);
2435 writer.write_batches(&[&batch]).unwrap();
2436 writer.finish().unwrap();
2437 }
2438
2439 assert_json_eq(
2440 &buf,
2441 r#"{"union":1,"float":1.0}
2442{"union":"a","float":null}
2443{"union":null,"float":3.4}
2444"#,
2445 );
2446 }
2447
2448 #[test]
2449 fn test_fallback_encoder_factory_array_implicit_nulls() {
2450 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2451
2452 let json_value: Value = {
2453 let mut buf = Vec::new();
2454 let mut writer = WriterBuilder::new()
2455 .with_encoder_factory(encoder_factory)
2456 .build::<_, JsonArray>(&mut buf);
2457 writer.write_batches(&[&batch]).unwrap();
2458 writer.finish().unwrap();
2459 serde_json::from_slice(&buf).unwrap()
2460 };
2461
2462 let expected = json!([
2463 {"union":1,"float":1.0},
2464 {"union":"a"},
2465 {"float":3.4,"union":null},
2466 ]);
2467
2468 assert_eq!(json_value, expected);
2469 }
2470
2471 #[test]
2472 fn test_fallback_encoder_factory_array_explicit_nulls() {
2473 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2474
2475 let json_value: Value = {
2476 let mut buf = Vec::new();
2477 let mut writer = WriterBuilder::new()
2478 .with_encoder_factory(encoder_factory)
2479 .with_explicit_nulls(true)
2480 .build::<_, JsonArray>(&mut buf);
2481 writer.write_batches(&[&batch]).unwrap();
2482 writer.finish().unwrap();
2483 serde_json::from_slice(&buf).unwrap()
2484 };
2485
2486 let expected = json!([
2487 {"union":1,"float":1.0},
2488 {"union":"a", "float": null},
2489 {"union":null,"float":3.4},
2490 ]);
2491
2492 assert_eq!(json_value, expected);
2493 }
2494
2495 #[test]
2496 fn test_default_encoder_byte_array() {
2497 struct IntArrayBinaryEncoder<B> {
2498 array: B,
2499 }
2500
2501 impl<'a, B> Encoder for IntArrayBinaryEncoder<B>
2502 where
2503 B: ArrayAccessor<Item = &'a [u8]>,
2504 {
2505 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2506 out.push(b'[');
2507 let child = self.array.value(idx);
2508 for (idx, byte) in child.iter().enumerate() {
2509 write!(out, "{byte}").unwrap();
2510 if idx < child.len() - 1 {
2511 out.push(b',');
2512 }
2513 }
2514 out.push(b']');
2515 }
2516 }
2517
2518 #[derive(Debug)]
2519 struct IntArayBinaryEncoderFactory;
2520
2521 impl EncoderFactory for IntArayBinaryEncoderFactory {
2522 fn make_default_encoder<'a>(
2523 &self,
2524 _field: &'a FieldRef,
2525 array: &'a dyn Array,
2526 _options: &'a EncoderOptions,
2527 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2528 match array.data_type() {
2529 DataType::Binary => {
2530 let array = array.as_binary::<i32>();
2531 let encoder = IntArrayBinaryEncoder { array };
2532 let array_encoder = Box::new(encoder) as Box<dyn Encoder + 'a>;
2533 let nulls = array.nulls().cloned();
2534 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2535 }
2536 _ => Ok(None),
2537 }
2538 }
2539 }
2540
2541 let binary_array = BinaryArray::from_opt_vec(vec![Some(b"a"), None, Some(b"b")]);
2542 let float_array = Float64Array::from(vec![Some(1.0), Some(2.3), None]);
2543 let fields = vec![
2544 Field::new("bytes", DataType::Binary, true),
2545 Field::new("float", DataType::Float64, true),
2546 ];
2547 let batch = RecordBatch::try_new(
2548 Arc::new(Schema::new(fields)),
2549 vec![
2550 Arc::new(binary_array) as Arc<dyn Array>,
2551 Arc::new(float_array) as Arc<dyn Array>,
2552 ],
2553 )
2554 .unwrap();
2555
2556 let json_value: Value = {
2557 let mut buf = Vec::new();
2558 let mut writer = WriterBuilder::new()
2559 .with_encoder_factory(Arc::new(IntArayBinaryEncoderFactory))
2560 .build::<_, JsonArray>(&mut buf);
2561 writer.write_batches(&[&batch]).unwrap();
2562 writer.finish().unwrap();
2563 serde_json::from_slice(&buf).unwrap()
2564 };
2565
2566 let expected = json!([
2567 {"bytes": [97], "float": 1.0},
2568 {"float": 2.3},
2569 {"bytes": [98]},
2570 ]);
2571
2572 assert_eq!(json_value, expected);
2573 }
2574
2575 #[test]
2576 fn test_encoder_factory_customize_dictionary() {
2577 struct PaddedInt32Encoder {
2582 array: Int32Array,
2583 }
2584
2585 impl Encoder for PaddedInt32Encoder {
2586 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2587 let value = self.array.value(idx);
2588 write!(out, "\"{value:0>8}\"").unwrap();
2589 }
2590 }
2591
2592 #[derive(Debug)]
2593 struct CustomEncoderFactory;
2594
2595 impl EncoderFactory for CustomEncoderFactory {
2596 fn make_default_encoder<'a>(
2597 &self,
2598 field: &'a FieldRef,
2599 array: &'a dyn Array,
2600 _options: &'a EncoderOptions,
2601 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2602 let padded = field
2607 .metadata()
2608 .get("padded")
2609 .map(|v| v == "true")
2610 .unwrap_or_default();
2611 match (array.data_type(), padded) {
2612 (DataType::Int32, true) => {
2613 let array = array.as_primitive::<Int32Type>();
2614 let nulls = array.nulls().cloned();
2615 let encoder = PaddedInt32Encoder {
2616 array: array.clone(),
2617 };
2618 let array_encoder = Box::new(encoder) as Box<dyn Encoder + 'a>;
2619 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2620 }
2621 _ => Ok(None),
2622 }
2623 }
2624 }
2625
2626 let to_json = |batch| {
2627 let mut buf = Vec::new();
2628 let mut writer = WriterBuilder::new()
2629 .with_encoder_factory(Arc::new(CustomEncoderFactory))
2630 .build::<_, JsonArray>(&mut buf);
2631 writer.write_batches(&[batch]).unwrap();
2632 writer.finish().unwrap();
2633 serde_json::from_slice::<Value>(&buf).unwrap()
2634 };
2635
2636 let array = Int32Array::from(vec![Some(1), None, Some(2)]);
2638 let field = Arc::new(Field::new("int", DataType::Int32, true).with_metadata(
2639 HashMap::from_iter(vec![("padded".to_string(), "true".to_string())]),
2640 ));
2641 let batch = RecordBatch::try_new(
2642 Arc::new(Schema::new(vec![field.clone()])),
2643 vec![Arc::new(array)],
2644 )
2645 .unwrap();
2646
2647 let json_value = to_json(&batch);
2648
2649 let expected = json!([
2650 {"int": "00000001"},
2651 {},
2652 {"int": "00000002"},
2653 ]);
2654
2655 assert_eq!(json_value, expected);
2656
2657 let mut array_builder = PrimitiveDictionaryBuilder::<UInt16Type, Int32Type>::new();
2659 array_builder.append_value(1);
2660 array_builder.append_null();
2661 array_builder.append_value(1);
2662 let array = array_builder.finish();
2663 let field = Field::new(
2664 "int",
2665 DataType::Dictionary(Box::new(DataType::UInt16), Box::new(DataType::Int32)),
2666 true,
2667 )
2668 .with_metadata(HashMap::from_iter(vec![(
2669 "padded".to_string(),
2670 "true".to_string(),
2671 )]));
2672 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![Arc::new(array)])
2673 .unwrap();
2674
2675 let json_value = to_json(&batch);
2676
2677 let expected = json!([
2678 {"int": "00000001"},
2679 {},
2680 {"int": "00000001"},
2681 ]);
2682
2683 assert_eq!(json_value, expected);
2684 }
2685
2686 #[test]
2687 fn test_write_run_end_encoded() {
2688 let run_ends = Int32Array::from(vec![2, 5, 6]);
2689 let values = StringArray::from(vec![Some("a"), Some("b"), None]);
2690 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2691
2692 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2693 "c1",
2694 ree.data_type().clone(),
2695 true,
2696 )]));
2697
2698 let batch = RecordBatch::try_new(schema, vec![Arc::new(ree)]).unwrap();
2699
2700 let mut buf = Vec::new();
2701 {
2702 let mut writer = LineDelimitedWriter::new(&mut buf);
2703 writer.write_batches(&[&batch]).unwrap();
2704 }
2705
2706 assert_json_eq(
2707 &buf,
2708 r#"{"c1":"a"}
2709{"c1":"a"}
2710{"c1":"b"}
2711{"c1":"b"}
2712{"c1":"b"}
2713{}
2714"#,
2715 );
2716 }
2717
2718 #[test]
2719 fn test_write_run_end_encoded_int_values() {
2720 let run_ends = Int32Array::from(vec![3, 5]);
2721 let values = Int32Array::from(vec![10, 20]);
2722 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2723
2724 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2725 "n",
2726 ree.data_type().clone(),
2727 true,
2728 )]));
2729
2730 let batch = RecordBatch::try_new(schema, vec![Arc::new(ree)]).unwrap();
2731
2732 let json_value: Value = {
2733 let mut buf = Vec::new();
2734 let mut writer = WriterBuilder::new().build::<_, JsonArray>(&mut buf);
2735 writer.write_batches(&[&batch]).unwrap();
2736 writer.finish().unwrap();
2737 serde_json::from_slice(&buf).unwrap()
2738 };
2739
2740 let expected = json!([
2741 {"n": 10},
2742 {"n": 10},
2743 {"n": 10},
2744 {"n": 20},
2745 {"n": 20},
2746 ]);
2747
2748 assert_eq!(json_value, expected);
2749 }
2750
2751 #[test]
2752 fn test_run_end_encoded_roundtrip() {
2753 let run_ends = Int32Array::from(vec![3, 5, 7]);
2754 let values = StringArray::from(vec![Some("a"), None, Some("b")]);
2755 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2756
2757 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2758 "c",
2759 ree.data_type().clone(),
2760 true,
2761 )]));
2762 let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(ree)]).unwrap();
2763
2764 let mut buf = Vec::new();
2765 {
2766 let mut writer = super::LineDelimitedWriter::new(&mut buf);
2767 writer.write_batches(&[&batch]).unwrap();
2768 }
2769
2770 let batches: Vec<RecordBatch> = ReaderBuilder::new(schema)
2771 .with_batch_size(1024)
2772 .build(std::io::Cursor::new(&buf))
2773 .unwrap()
2774 .collect::<Result<Vec<_>, _>>()
2775 .unwrap();
2776 assert_eq!(batches.len(), 1);
2777
2778 let col = batches[0].column(0);
2779 let run_array = col.as_run::<Int32Type>();
2780
2781 assert_eq!(run_array.len(), 7);
2782 assert_eq!(run_array.run_ends().values(), &[3, 5, 7]);
2783
2784 let values = run_array.values().as_string::<i32>();
2785 assert_eq!(values.len(), 3);
2786 assert_eq!(values.value(0), "a");
2787 assert!(values.is_null(1));
2788 assert_eq!(values.value(2), "b");
2789 }
2790}