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 #[cfg_attr(miri, ignore)] fn test_write_multi_batches() {
1526 let test_file = "test/data/basic.json";
1527
1528 let schema = SchemaRef::new(Schema::new(vec![
1529 Field::new("a", DataType::Int64, true),
1530 Field::new("b", DataType::Float64, true),
1531 Field::new("c", DataType::Boolean, true),
1532 Field::new("d", DataType::Utf8, true),
1533 Field::new("e", DataType::Utf8, true),
1534 Field::new("f", DataType::Utf8, true),
1535 Field::new("g", DataType::Timestamp(TimeUnit::Millisecond, None), true),
1536 Field::new("h", DataType::Float16, true),
1537 ]));
1538
1539 let mut reader = ReaderBuilder::new(schema.clone())
1540 .build(BufReader::new(File::open(test_file).unwrap()))
1541 .unwrap();
1542 let batch = reader.next().unwrap().unwrap();
1543
1544 let batches = [&RecordBatch::new_empty(schema), &batch, &batch];
1546
1547 let mut buf = Vec::new();
1548 {
1549 let mut writer = LineDelimitedWriter::new(&mut buf);
1550 writer.write_batches(&batches).unwrap();
1551 }
1552
1553 let result = str::from_utf8(&buf).unwrap();
1554 let expected = read_to_string(test_file).unwrap();
1555 let expected = format!("{expected}\n{expected}");
1557 for (r, e) in result.lines().zip(expected.lines()) {
1558 let mut expected_json = serde_json::from_str::<Value>(e).unwrap();
1559 if let Value::Object(obj) = expected_json {
1561 expected_json =
1562 Value::Object(obj.into_iter().filter(|(_, v)| *v != Value::Null).collect());
1563 }
1564 assert_eq!(serde_json::from_str::<Value>(r).unwrap(), expected_json,);
1565 }
1566 }
1567
1568 #[test]
1569 fn test_writer_explicit_nulls() -> Result<(), ArrowError> {
1570 fn nested_list() -> (Arc<ListArray>, Arc<Field>) {
1571 let array = Arc::new(ListArray::from_iter_primitive::<Int32Type, _, _>(vec![
1572 Some(vec![None, None, None]),
1573 Some(vec![Some(1), Some(2), Some(3)]),
1574 None,
1575 Some(vec![None, None, None]),
1576 ]));
1577 let field = Arc::new(Field::new("list", array.data_type().clone(), true));
1578 (array, field)
1580 }
1581
1582 fn nested_dict() -> (Arc<DictionaryArray<Int32Type>>, Arc<Field>) {
1583 let array = Arc::new(DictionaryArray::from_iter(vec![
1584 Some("cupcakes"),
1585 None,
1586 Some("bear"),
1587 Some("kuma"),
1588 ]));
1589 let field = Arc::new(Field::new("dict", array.data_type().clone(), true));
1590 (array, field)
1592 }
1593
1594 fn nested_map() -> (Arc<MapArray>, Arc<Field>) {
1595 let string_builder = StringBuilder::new();
1596 let int_builder = Int64Builder::new();
1597 let mut builder = MapBuilder::new(None, string_builder, int_builder);
1598
1599 builder.keys().append_value("foo");
1601 builder.values().append_value(10);
1602 builder.append(true).unwrap();
1603
1604 builder.append(false).unwrap();
1605
1606 builder.append(true).unwrap();
1607
1608 builder.keys().append_value("bar");
1609 builder.values().append_value(20);
1610 builder.keys().append_value("baz");
1611 builder.values().append_value(30);
1612 builder.keys().append_value("qux");
1613 builder.values().append_value(40);
1614 builder.append(true).unwrap();
1615
1616 let array = Arc::new(builder.finish());
1617 let field = Arc::new(Field::new("map", array.data_type().clone(), true));
1618 (array, field)
1619 }
1620
1621 fn root_list() -> (Arc<ListArray>, Field) {
1622 let struct_array = StructArray::from(vec![
1623 (
1624 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1625 Arc::new(StringArray::from(vec![Some("a"), Some("b"), None, None])) as ArrayRef,
1626 ),
1627 (
1628 Arc::new(Field::new("int32", DataType::Int32, true)),
1629 Arc::new(Int32Array::from(vec![Some(1), None, Some(5), None])) as ArrayRef,
1630 ),
1631 ]);
1632
1633 let values_field =
1634 Arc::new(Field::new("struct", struct_array.data_type().clone(), true));
1635 let field = Field::new_list("list", values_field.as_ref().clone(), true);
1636
1637 let array = Arc::new(ListArray::new(
1639 values_field,
1640 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 2, 3, 3])),
1641 Arc::new(struct_array),
1642 Some(NullBuffer::from(vec![true, false, true, false])),
1643 ));
1644 (array, field)
1645 }
1646
1647 let (nested_list_array, nested_list_field) = nested_list();
1648 let (nested_dict_array, nested_dict_field) = nested_dict();
1649 let (nested_map_array, nested_map_field) = nested_map();
1650 let (root_list_array, root_list_field) = root_list();
1651
1652 let schema = Schema::new(vec![
1653 Field::new("date", DataType::Date32, true),
1654 Field::new("null", DataType::Null, true),
1655 Field::new_struct(
1656 "struct",
1657 vec![
1658 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1659 nested_list_field.clone(),
1660 nested_dict_field.clone(),
1661 nested_map_field.clone(),
1662 ],
1663 true,
1664 ),
1665 root_list_field,
1666 ]);
1667
1668 let arr_date32 = Date32Array::from(vec![Some(0), None, Some(1), None]);
1669 let arr_null = NullArray::new(4);
1670 let arr_struct = StructArray::from(vec![
1671 (
1673 Arc::new(Field::new("utf8", DataType::Utf8, true)),
1674 Arc::new(StringArray::from(vec![Some("a"), None, None, Some("b")])) as ArrayRef,
1675 ),
1676 (nested_list_field, nested_list_array as ArrayRef),
1678 (nested_dict_field, nested_dict_array as ArrayRef),
1680 (nested_map_field, nested_map_array as ArrayRef),
1682 ]);
1683
1684 let batch = RecordBatch::try_new(
1685 Arc::new(schema),
1686 vec![
1687 Arc::new(arr_date32),
1689 Arc::new(arr_null),
1691 Arc::new(arr_struct),
1692 root_list_array,
1694 ],
1695 )?;
1696
1697 let mut buf = Vec::new();
1698 {
1699 let mut writer = WriterBuilder::new()
1700 .with_explicit_nulls(true)
1701 .build::<_, JsonArray>(&mut buf);
1702 writer.write_batches(&[&batch])?;
1703 writer.finish()?;
1704 }
1705
1706 let actual = serde_json::from_slice::<Vec<Value>>(&buf).unwrap();
1707 let expected = serde_json::from_value::<Vec<Value>>(json!([
1708 {
1709 "date": "1970-01-01",
1710 "list": [
1711 {
1712 "int32": 1,
1713 "utf8": "a"
1714 },
1715 {
1716 "int32": null,
1717 "utf8": "b"
1718 }
1719 ],
1720 "null": null,
1721 "struct": {
1722 "dict": "cupcakes",
1723 "list": [
1724 null,
1725 null,
1726 null
1727 ],
1728 "map": {
1729 "foo": 10
1730 },
1731 "utf8": "a"
1732 }
1733 },
1734 {
1735 "date": null,
1736 "list": null,
1737 "null": null,
1738 "struct": {
1739 "dict": null,
1740 "list": [
1741 1,
1742 2,
1743 3
1744 ],
1745 "map": null,
1746 "utf8": null
1747 }
1748 },
1749 {
1750 "date": "1970-01-02",
1751 "list": [
1752 {
1753 "int32": 5,
1754 "utf8": null
1755 }
1756 ],
1757 "null": null,
1758 "struct": {
1759 "dict": "bear",
1760 "list": null,
1761 "map": {},
1762 "utf8": null
1763 }
1764 },
1765 {
1766 "date": null,
1767 "list": null,
1768 "null": null,
1769 "struct": {
1770 "dict": "kuma",
1771 "list": [
1772 null,
1773 null,
1774 null
1775 ],
1776 "map": {
1777 "bar": 20,
1778 "baz": 30,
1779 "qux": 40
1780 },
1781 "utf8": "b"
1782 }
1783 }
1784 ]))
1785 .unwrap();
1786
1787 assert_eq!(actual, expected);
1788
1789 Ok(())
1790 }
1791
1792 fn build_array_binary<O: OffsetSizeTrait>(values: &[Option<&[u8]>]) -> RecordBatch {
1793 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1794 "bytes",
1795 GenericBinaryType::<O>::DATA_TYPE,
1796 true,
1797 )]));
1798 let mut builder = GenericByteBuilder::<GenericBinaryType<O>>::new();
1799 for value in values {
1800 match value {
1801 Some(v) => builder.append_value(v),
1802 None => builder.append_null(),
1803 }
1804 }
1805 let array = Arc::new(builder.finish()) as ArrayRef;
1806 RecordBatch::try_new(schema, vec![array]).unwrap()
1807 }
1808
1809 fn build_array_binary_view(values: &[Option<&[u8]>]) -> RecordBatch {
1810 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1811 "bytes",
1812 DataType::BinaryView,
1813 true,
1814 )]));
1815 let mut builder = BinaryViewBuilder::new();
1816 for value in values {
1817 match value {
1818 Some(v) => builder.append_value(v),
1819 None => builder.append_null(),
1820 }
1821 }
1822 let array = Arc::new(builder.finish()) as ArrayRef;
1823 RecordBatch::try_new(schema, vec![array]).unwrap()
1824 }
1825
1826 fn assert_binary_json(batch: &RecordBatch) {
1827 {
1829 let mut buf = Vec::new();
1830 let json_value: Value = {
1831 let mut writer = WriterBuilder::new()
1832 .with_explicit_nulls(true)
1833 .build::<_, JsonArray>(&mut buf);
1834 writer.write(batch).unwrap();
1835 writer.close().unwrap();
1836 serde_json::from_slice(&buf).unwrap()
1837 };
1838
1839 assert_eq!(
1840 json!([
1841 {
1842 "bytes": "4e656420466c616e64657273"
1843 },
1844 {
1845 "bytes": null },
1847 {
1848 "bytes": "54726f79204d63436c757265"
1849 }
1850 ]),
1851 json_value,
1852 );
1853 }
1854
1855 {
1857 let mut buf = Vec::new();
1858 let json_value: Value = {
1859 let mut writer = ArrayWriter::new(&mut buf);
1862 writer.write(batch).unwrap();
1863 writer.close().unwrap();
1864 serde_json::from_slice(&buf).unwrap()
1865 };
1866
1867 assert_eq!(
1868 json!([
1869 { "bytes": "4e656420466c616e64657273" },
1870 {},
1871 { "bytes": "54726f79204d63436c757265" }
1872 ]),
1873 json_value
1874 );
1875 }
1876 }
1877
1878 #[test]
1879 fn test_writer_binary() {
1880 let values: [Option<&[u8]>; 3] = [
1881 Some(b"Ned Flanders" as &[u8]),
1882 None,
1883 Some(b"Troy McClure" as &[u8]),
1884 ];
1885 {
1887 let batch = build_array_binary::<i32>(&values);
1888 assert_binary_json(&batch);
1889 }
1890 {
1892 let batch = build_array_binary::<i64>(&values);
1893 assert_binary_json(&batch);
1894 }
1895 {
1896 let batch = build_array_binary_view(&values);
1897 assert_binary_json(&batch);
1898 }
1899 }
1900
1901 #[test]
1902 fn test_writer_fixed_size_binary() {
1903 let size = 11;
1905 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1906 "bytes",
1907 DataType::FixedSizeBinary(size),
1908 true,
1909 )]));
1910
1911 let mut builder = FixedSizeBinaryBuilder::new(size);
1913 let values = [Some(b"hello world"), None, Some(b"summer rain")];
1914 for value in values {
1915 match value {
1916 Some(v) => builder.append_value(v).unwrap(),
1917 None => builder.append_null(),
1918 }
1919 }
1920 let array = Arc::new(builder.finish()) as ArrayRef;
1921 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
1922
1923 {
1925 let mut buf = Vec::new();
1926 let json_value: Value = {
1927 let mut writer = WriterBuilder::new()
1928 .with_explicit_nulls(true)
1929 .build::<_, JsonArray>(&mut buf);
1930 writer.write(&batch).unwrap();
1931 writer.close().unwrap();
1932 serde_json::from_slice(&buf).unwrap()
1933 };
1934
1935 assert_eq!(
1936 json!([
1937 {
1938 "bytes": "68656c6c6f20776f726c64"
1939 },
1940 {
1941 "bytes": null },
1943 {
1944 "bytes": "73756d6d6572207261696e"
1945 }
1946 ]),
1947 json_value,
1948 );
1949 }
1950 {
1952 let mut buf = Vec::new();
1953 let json_value: Value = {
1954 let mut writer = ArrayWriter::new(&mut buf);
1957 writer.write(&batch).unwrap();
1958 writer.close().unwrap();
1959 serde_json::from_slice(&buf).unwrap()
1960 };
1961
1962 assert_eq!(
1963 json!([
1964 {
1965 "bytes": "68656c6c6f20776f726c64"
1966 },
1967 {}, {
1969 "bytes": "73756d6d6572207261696e"
1970 }
1971 ]),
1972 json_value,
1973 );
1974 }
1975 }
1976
1977 #[test]
1978 fn test_writer_fixed_size_list() {
1979 let size = 3;
1980 let field = FieldRef::new(Field::new_list_field(DataType::Int32, true));
1981 let schema = SchemaRef::new(Schema::new(vec![Field::new(
1982 "list",
1983 DataType::FixedSizeList(field, size),
1984 true,
1985 )]));
1986
1987 let values_builder = Int32Builder::new();
1988 let mut list_builder = FixedSizeListBuilder::new(values_builder, size);
1989 let lists = [
1990 Some([Some(1), Some(2), None]),
1991 Some([Some(3), None, Some(4)]),
1992 Some([None, Some(5), Some(6)]),
1993 None,
1994 ];
1995 for list in lists {
1996 match list {
1997 Some(l) => {
1998 for value in l {
1999 match value {
2000 Some(v) => list_builder.values().append_value(v),
2001 None => list_builder.values().append_null(),
2002 }
2003 }
2004 list_builder.append(true);
2005 }
2006 None => {
2007 for _ in 0..size {
2008 list_builder.values().append_null();
2009 }
2010 list_builder.append(false);
2011 }
2012 }
2013 }
2014 let array = Arc::new(list_builder.finish()) as ArrayRef;
2015 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
2016
2017 {
2019 let json_value: Value = {
2020 let mut buf = Vec::new();
2021 let mut writer = WriterBuilder::new()
2022 .with_explicit_nulls(true)
2023 .build::<_, JsonArray>(&mut buf);
2024 writer.write(&batch).unwrap();
2025 writer.close().unwrap();
2026 serde_json::from_slice(&buf).unwrap()
2027 };
2028 assert_eq!(
2029 json!([
2030 {"list": [1, 2, null]},
2031 {"list": [3, null, 4]},
2032 {"list": [null, 5, 6]},
2033 {"list": null},
2034 ]),
2035 json_value
2036 );
2037 }
2038 {
2040 let json_value: Value = {
2041 let mut buf = Vec::new();
2042 let mut writer = ArrayWriter::new(&mut buf);
2043 writer.write(&batch).unwrap();
2044 writer.close().unwrap();
2045 serde_json::from_slice(&buf).unwrap()
2046 };
2047 assert_eq!(
2048 json!([
2049 {"list": [1, 2, null]},
2050 {"list": [3, null, 4]},
2051 {"list": [null, 5, 6]},
2052 {}, ]),
2054 json_value
2055 );
2056 }
2057 }
2058
2059 #[test]
2060 fn test_writer_null_dict() {
2061 let keys = Int32Array::from_iter(vec![Some(0), None, Some(1)]);
2062 let values = Arc::new(StringArray::from_iter(vec![Some("a"), None]));
2063 let dict = DictionaryArray::new(keys, values);
2064
2065 let schema = SchemaRef::new(Schema::new(vec![Field::new(
2066 "my_dict",
2067 DataType::Dictionary(DataType::Int32.into(), DataType::Utf8.into()),
2068 true,
2069 )]));
2070
2071 let array = Arc::new(dict) as ArrayRef;
2072 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
2073
2074 let mut json = Vec::new();
2075 let write_builder = WriterBuilder::new().with_explicit_nulls(true);
2076 let mut writer = write_builder.build::<_, JsonArray>(&mut json);
2077 writer.write(&batch).unwrap();
2078 writer.close().unwrap();
2079
2080 let json_str = str::from_utf8(&json).unwrap();
2081 assert_eq!(
2082 json_str,
2083 r#"[{"my_dict":"a"},{"my_dict":null},{"my_dict":""}]"#
2084 )
2085 }
2086
2087 #[test]
2088 fn test_decimal32_encoder() {
2089 let array = Decimal32Array::from_iter_values([1234, 5678, 9012])
2090 .with_precision_and_scale(8, 2)
2091 .unwrap();
2092 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2093 let schema = Schema::new(vec![field]);
2094 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2095
2096 let mut buf = Vec::new();
2097 {
2098 let mut writer = LineDelimitedWriter::new(&mut buf);
2099 writer.write_batches(&[&batch]).unwrap();
2100 }
2101
2102 assert_json_eq(
2103 &buf,
2104 r#"{"decimal":12.34}
2105{"decimal":56.78}
2106{"decimal":90.12}
2107"#,
2108 );
2109 }
2110
2111 #[test]
2112 fn test_decimal64_encoder() {
2113 let array = Decimal64Array::from_iter_values([1234, 5678, 9012])
2114 .with_precision_and_scale(10, 2)
2115 .unwrap();
2116 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2117 let schema = Schema::new(vec![field]);
2118 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2119
2120 let mut buf = Vec::new();
2121 {
2122 let mut writer = LineDelimitedWriter::new(&mut buf);
2123 writer.write_batches(&[&batch]).unwrap();
2124 }
2125
2126 assert_json_eq(
2127 &buf,
2128 r#"{"decimal":12.34}
2129{"decimal":56.78}
2130{"decimal":90.12}
2131"#,
2132 );
2133 }
2134
2135 #[test]
2136 fn test_decimal128_encoder() {
2137 let array = Decimal128Array::from_iter_values([1234, 5678, 9012])
2138 .with_precision_and_scale(10, 2)
2139 .unwrap();
2140 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2141 let schema = Schema::new(vec![field]);
2142 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2143
2144 let mut buf = Vec::new();
2145 {
2146 let mut writer = LineDelimitedWriter::new(&mut buf);
2147 writer.write_batches(&[&batch]).unwrap();
2148 }
2149
2150 assert_json_eq(
2151 &buf,
2152 r#"{"decimal":12.34}
2153{"decimal":56.78}
2154{"decimal":90.12}
2155"#,
2156 );
2157 }
2158
2159 #[test]
2160 fn test_decimal256_encoder() {
2161 let array = Decimal256Array::from_iter_values([
2162 i256::from(123400),
2163 i256::from(567800),
2164 i256::from(901200),
2165 ])
2166 .with_precision_and_scale(10, 4)
2167 .unwrap();
2168 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2169 let schema = Schema::new(vec![field]);
2170 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2171
2172 let mut buf = Vec::new();
2173 {
2174 let mut writer = LineDelimitedWriter::new(&mut buf);
2175 writer.write_batches(&[&batch]).unwrap();
2176 }
2177
2178 assert_json_eq(
2179 &buf,
2180 r#"{"decimal":12.3400}
2181{"decimal":56.7800}
2182{"decimal":90.1200}
2183"#,
2184 );
2185 }
2186
2187 #[test]
2188 fn test_decimal_encoder_with_nulls() {
2189 let array = Decimal128Array::from_iter([Some(1234), None, Some(5678)])
2190 .with_precision_and_scale(10, 2)
2191 .unwrap();
2192 let field = Arc::new(Field::new("decimal", array.data_type().clone(), true));
2193 let schema = Schema::new(vec![field]);
2194 let batch = RecordBatch::try_new(Arc::new(schema), vec![Arc::new(array)]).unwrap();
2195
2196 let mut buf = Vec::new();
2197 {
2198 let mut writer = LineDelimitedWriter::new(&mut buf);
2199 writer.write_batches(&[&batch]).unwrap();
2200 }
2201
2202 assert_json_eq(
2203 &buf,
2204 r#"{"decimal":12.34}
2205{}
2206{"decimal":56.78}
2207"#,
2208 );
2209 }
2210
2211 #[test]
2212 fn write_structs_as_list() {
2213 let schema = Schema::new(vec![
2214 Field::new(
2215 "c1",
2216 DataType::Struct(Fields::from(vec![
2217 Field::new("c11", DataType::Int32, true),
2218 Field::new(
2219 "c12",
2220 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
2221 false,
2222 ),
2223 ])),
2224 false,
2225 ),
2226 Field::new("c2", DataType::Utf8, false),
2227 ]);
2228
2229 let c1 = StructArray::from(vec![
2230 (
2231 Arc::new(Field::new("c11", DataType::Int32, true)),
2232 Arc::new(Int32Array::from(vec![Some(1), None, Some(5)])) as ArrayRef,
2233 ),
2234 (
2235 Arc::new(Field::new(
2236 "c12",
2237 DataType::Struct(vec![Field::new("c121", DataType::Utf8, false)].into()),
2238 false,
2239 )),
2240 Arc::new(StructArray::from(vec![(
2241 Arc::new(Field::new("c121", DataType::Utf8, false)),
2242 Arc::new(StringArray::from(vec![Some("e"), Some("f"), Some("g")])) as ArrayRef,
2243 )])) as ArrayRef,
2244 ),
2245 ]);
2246 let c2 = StringArray::from(vec![Some("a"), Some("b"), Some("c")]);
2247
2248 let batch =
2249 RecordBatch::try_new(Arc::new(schema), vec![Arc::new(c1), Arc::new(c2)]).unwrap();
2250
2251 let expected = r#"[[1,["e"]],"a"]
2252[[null,["f"]],"b"]
2253[[5,["g"]],"c"]
2254"#;
2255
2256 let mut buf = Vec::new();
2257 {
2258 let builder = WriterBuilder::new()
2259 .with_explicit_nulls(true)
2260 .with_struct_mode(StructMode::ListOnly);
2261 let mut writer = builder.build::<_, LineDelimited>(&mut buf);
2262 writer.write_batches(&[&batch]).unwrap();
2263 }
2264 assert_json_eq(&buf, expected);
2265
2266 let mut buf = Vec::new();
2267 {
2268 let builder = WriterBuilder::new()
2269 .with_explicit_nulls(false)
2270 .with_struct_mode(StructMode::ListOnly);
2271 let mut writer = builder.build::<_, LineDelimited>(&mut buf);
2272 writer.write_batches(&[&batch]).unwrap();
2273 }
2274 assert_json_eq(&buf, expected);
2275 }
2276
2277 fn make_fallback_encoder_test_data() -> (RecordBatch, Arc<dyn EncoderFactory>) {
2278 #[derive(Debug)]
2281 enum UnionValue {
2282 Int32(i32),
2283 String(String),
2284 }
2285
2286 #[derive(Debug)]
2287 struct UnionEncoder {
2288 array: Vec<Option<UnionValue>>,
2289 }
2290
2291 impl Encoder for UnionEncoder {
2292 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2293 match &self.array[idx] {
2294 None => out.extend_from_slice(b"null"),
2295 Some(UnionValue::Int32(v)) => out.extend_from_slice(v.to_string().as_bytes()),
2296 Some(UnionValue::String(v)) => {
2297 out.extend_from_slice(format!("\"{v}\"").as_bytes())
2298 }
2299 }
2300 }
2301 }
2302
2303 #[derive(Debug)]
2304 struct UnionEncoderFactory;
2305
2306 impl EncoderFactory for UnionEncoderFactory {
2307 fn make_default_encoder<'a>(
2308 &self,
2309 _field: &'a FieldRef,
2310 array: &'a dyn Array,
2311 _options: &'a EncoderOptions,
2312 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2313 let data_type = array.data_type();
2314 let fields = match data_type {
2315 DataType::Union(fields, UnionMode::Sparse) => fields,
2316 _ => return Ok(None),
2317 };
2318 let fields = fields.iter().map(|(_, f)| f).collect::<Vec<_>>();
2320 for f in fields.iter() {
2321 match f.data_type() {
2322 DataType::Null => {}
2323 DataType::Int32 => {}
2324 DataType::Utf8 => {}
2325 _ => return Ok(None),
2326 }
2327 }
2328 let (_, type_ids, _, buffers) = array.as_union().clone().into_parts();
2329 let mut values = Vec::with_capacity(type_ids.len());
2330 for idx in 0..type_ids.len() {
2331 let type_id = type_ids[idx];
2332 let field = &fields[type_id as usize];
2333 let value = match field.data_type() {
2334 DataType::Null => None,
2335 DataType::Int32 => Some(UnionValue::Int32(
2336 buffers[type_id as usize]
2337 .as_primitive::<Int32Type>()
2338 .value(idx),
2339 )),
2340 DataType::Utf8 => Some(UnionValue::String(
2341 buffers[type_id as usize]
2342 .as_string::<i32>()
2343 .value(idx)
2344 .to_string(),
2345 )),
2346 _ => unreachable!(),
2347 };
2348 values.push(value);
2349 }
2350 let array_encoder =
2351 Box::new(UnionEncoder { array: values }) as Box<dyn Encoder + 'a>;
2352 let nulls = array.nulls().cloned();
2353 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2354 }
2355 }
2356
2357 let int_array = Int32Array::from(vec![Some(1), None, None]);
2358 let string_array = StringArray::from(vec![None, Some("a"), None]);
2359 let null_array = NullArray::new(3);
2360 let type_ids = [0_i8, 1, 2].into_iter().collect::<ScalarBuffer<i8>>();
2361
2362 let union_fields = [
2363 (0, Arc::new(Field::new("A", DataType::Int32, false))),
2364 (1, Arc::new(Field::new("B", DataType::Utf8, false))),
2365 (2, Arc::new(Field::new("C", DataType::Null, false))),
2366 ]
2367 .into_iter()
2368 .collect::<UnionFields>();
2369
2370 let children = vec![
2371 Arc::new(int_array) as Arc<dyn Array>,
2372 Arc::new(string_array),
2373 Arc::new(null_array),
2374 ];
2375
2376 let array = UnionArray::try_new(union_fields.clone(), type_ids, None, children).unwrap();
2377
2378 let float_array = Float64Array::from(vec![Some(1.0), None, Some(3.4)]);
2379
2380 let fields = vec![
2381 Field::new(
2382 "union",
2383 DataType::Union(union_fields, UnionMode::Sparse),
2384 true,
2385 ),
2386 Field::new("float", DataType::Float64, true),
2387 ];
2388
2389 let batch = RecordBatch::try_new(
2390 Arc::new(Schema::new(fields)),
2391 vec![
2392 Arc::new(array) as Arc<dyn Array>,
2393 Arc::new(float_array) as Arc<dyn Array>,
2394 ],
2395 )
2396 .unwrap();
2397
2398 (batch, Arc::new(UnionEncoderFactory))
2399 }
2400
2401 #[test]
2402 fn test_fallback_encoder_factory_line_delimited_implicit_nulls() {
2403 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2404
2405 let mut buf = Vec::new();
2406 {
2407 let mut writer = WriterBuilder::new()
2408 .with_encoder_factory(encoder_factory)
2409 .with_explicit_nulls(false)
2410 .build::<_, LineDelimited>(&mut buf);
2411 writer.write_batches(&[&batch]).unwrap();
2412 writer.finish().unwrap();
2413 }
2414
2415 println!("{}", str::from_utf8(&buf).unwrap());
2416
2417 assert_json_eq(
2418 &buf,
2419 r#"{"union":1,"float":1.0}
2420{"union":"a"}
2421{"union":null,"float":3.4}
2422"#,
2423 );
2424 }
2425
2426 #[test]
2427 fn test_fallback_encoder_factory_line_delimited_explicit_nulls() {
2428 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2429
2430 let mut buf = Vec::new();
2431 {
2432 let mut writer = WriterBuilder::new()
2433 .with_encoder_factory(encoder_factory)
2434 .with_explicit_nulls(true)
2435 .build::<_, LineDelimited>(&mut buf);
2436 writer.write_batches(&[&batch]).unwrap();
2437 writer.finish().unwrap();
2438 }
2439
2440 assert_json_eq(
2441 &buf,
2442 r#"{"union":1,"float":1.0}
2443{"union":"a","float":null}
2444{"union":null,"float":3.4}
2445"#,
2446 );
2447 }
2448
2449 #[test]
2450 fn test_fallback_encoder_factory_array_implicit_nulls() {
2451 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2452
2453 let json_value: Value = {
2454 let mut buf = Vec::new();
2455 let mut writer = WriterBuilder::new()
2456 .with_encoder_factory(encoder_factory)
2457 .build::<_, JsonArray>(&mut buf);
2458 writer.write_batches(&[&batch]).unwrap();
2459 writer.finish().unwrap();
2460 serde_json::from_slice(&buf).unwrap()
2461 };
2462
2463 let expected = json!([
2464 {"union":1,"float":1.0},
2465 {"union":"a"},
2466 {"float":3.4,"union":null},
2467 ]);
2468
2469 assert_eq!(json_value, expected);
2470 }
2471
2472 #[test]
2473 fn test_fallback_encoder_factory_array_explicit_nulls() {
2474 let (batch, encoder_factory) = make_fallback_encoder_test_data();
2475
2476 let json_value: Value = {
2477 let mut buf = Vec::new();
2478 let mut writer = WriterBuilder::new()
2479 .with_encoder_factory(encoder_factory)
2480 .with_explicit_nulls(true)
2481 .build::<_, JsonArray>(&mut buf);
2482 writer.write_batches(&[&batch]).unwrap();
2483 writer.finish().unwrap();
2484 serde_json::from_slice(&buf).unwrap()
2485 };
2486
2487 let expected = json!([
2488 {"union":1,"float":1.0},
2489 {"union":"a", "float": null},
2490 {"union":null,"float":3.4},
2491 ]);
2492
2493 assert_eq!(json_value, expected);
2494 }
2495
2496 #[test]
2497 fn test_default_encoder_byte_array() {
2498 struct IntArrayBinaryEncoder<B> {
2499 array: B,
2500 }
2501
2502 impl<'a, B> Encoder for IntArrayBinaryEncoder<B>
2503 where
2504 B: ArrayAccessor<Item = &'a [u8]>,
2505 {
2506 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2507 out.push(b'[');
2508 let child = self.array.value(idx);
2509 for (idx, byte) in child.iter().enumerate() {
2510 write!(out, "{byte}").unwrap();
2511 if idx < child.len() - 1 {
2512 out.push(b',');
2513 }
2514 }
2515 out.push(b']');
2516 }
2517 }
2518
2519 #[derive(Debug)]
2520 struct IntArayBinaryEncoderFactory;
2521
2522 impl EncoderFactory for IntArayBinaryEncoderFactory {
2523 fn make_default_encoder<'a>(
2524 &self,
2525 _field: &'a FieldRef,
2526 array: &'a dyn Array,
2527 _options: &'a EncoderOptions,
2528 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2529 match array.data_type() {
2530 DataType::Binary => {
2531 let array = array.as_binary::<i32>();
2532 let encoder = IntArrayBinaryEncoder { array };
2533 let array_encoder = Box::new(encoder) as Box<dyn Encoder + 'a>;
2534 let nulls = array.nulls().cloned();
2535 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2536 }
2537 _ => Ok(None),
2538 }
2539 }
2540 }
2541
2542 let binary_array = BinaryArray::from_opt_vec(vec![Some(b"a"), None, Some(b"b")]);
2543 let float_array = Float64Array::from(vec![Some(1.0), Some(2.3), None]);
2544 let fields = vec![
2545 Field::new("bytes", DataType::Binary, true),
2546 Field::new("float", DataType::Float64, true),
2547 ];
2548 let batch = RecordBatch::try_new(
2549 Arc::new(Schema::new(fields)),
2550 vec![
2551 Arc::new(binary_array) as Arc<dyn Array>,
2552 Arc::new(float_array) as Arc<dyn Array>,
2553 ],
2554 )
2555 .unwrap();
2556
2557 let json_value: Value = {
2558 let mut buf = Vec::new();
2559 let mut writer = WriterBuilder::new()
2560 .with_encoder_factory(Arc::new(IntArayBinaryEncoderFactory))
2561 .build::<_, JsonArray>(&mut buf);
2562 writer.write_batches(&[&batch]).unwrap();
2563 writer.finish().unwrap();
2564 serde_json::from_slice(&buf).unwrap()
2565 };
2566
2567 let expected = json!([
2568 {"bytes": [97], "float": 1.0},
2569 {"float": 2.3},
2570 {"bytes": [98]},
2571 ]);
2572
2573 assert_eq!(json_value, expected);
2574 }
2575
2576 #[test]
2577 fn test_encoder_factory_customize_dictionary() {
2578 struct PaddedInt32Encoder {
2583 array: Int32Array,
2584 }
2585
2586 impl Encoder for PaddedInt32Encoder {
2587 fn encode(&mut self, idx: usize, out: &mut Vec<u8>) {
2588 let value = self.array.value(idx);
2589 write!(out, "\"{value:0>8}\"").unwrap();
2590 }
2591 }
2592
2593 #[derive(Debug)]
2594 struct CustomEncoderFactory;
2595
2596 impl EncoderFactory for CustomEncoderFactory {
2597 fn make_default_encoder<'a>(
2598 &self,
2599 field: &'a FieldRef,
2600 array: &'a dyn Array,
2601 _options: &'a EncoderOptions,
2602 ) -> Result<Option<NullableEncoder<'a>>, ArrowError> {
2603 let padded = field.metadata().get("padded").is_some_and(|v| v == "true");
2608 match (array.data_type(), padded) {
2609 (DataType::Int32, true) => {
2610 let array = array.as_primitive::<Int32Type>();
2611 let nulls = array.nulls().cloned();
2612 let encoder = PaddedInt32Encoder {
2613 array: array.clone(),
2614 };
2615 let array_encoder = Box::new(encoder) as Box<dyn Encoder + 'a>;
2616 Ok(Some(NullableEncoder::new(array_encoder, nulls)))
2617 }
2618 _ => Ok(None),
2619 }
2620 }
2621 }
2622
2623 let to_json = |batch| {
2624 let mut buf = Vec::new();
2625 let mut writer = WriterBuilder::new()
2626 .with_encoder_factory(Arc::new(CustomEncoderFactory))
2627 .build::<_, JsonArray>(&mut buf);
2628 writer.write_batches(&[batch]).unwrap();
2629 writer.finish().unwrap();
2630 serde_json::from_slice::<Value>(&buf).unwrap()
2631 };
2632
2633 let array = Int32Array::from(vec![Some(1), None, Some(2)]);
2635 let field = Arc::new(Field::new("int", DataType::Int32, true).with_metadata(
2636 HashMap::from_iter(vec![("padded".to_string(), "true".to_string())]),
2637 ));
2638 let batch = RecordBatch::try_new(
2639 Arc::new(Schema::new(vec![field.clone()])),
2640 vec![Arc::new(array)],
2641 )
2642 .unwrap();
2643
2644 let json_value = to_json(&batch);
2645
2646 let expected = json!([
2647 {"int": "00000001"},
2648 {},
2649 {"int": "00000002"},
2650 ]);
2651
2652 assert_eq!(json_value, expected);
2653
2654 let mut array_builder = PrimitiveDictionaryBuilder::<UInt16Type, Int32Type>::new();
2656 array_builder.append_value(1);
2657 array_builder.append_null();
2658 array_builder.append_value(1);
2659 let array = array_builder.finish();
2660 let field = Field::new(
2661 "int",
2662 DataType::Dictionary(Box::new(DataType::UInt16), Box::new(DataType::Int32)),
2663 true,
2664 )
2665 .with_metadata(HashMap::from_iter(vec![(
2666 "padded".to_string(),
2667 "true".to_string(),
2668 )]));
2669 let batch = RecordBatch::try_new(Arc::new(Schema::new(vec![field])), vec![Arc::new(array)])
2670 .unwrap();
2671
2672 let json_value = to_json(&batch);
2673
2674 let expected = json!([
2675 {"int": "00000001"},
2676 {},
2677 {"int": "00000001"},
2678 ]);
2679
2680 assert_eq!(json_value, expected);
2681 }
2682
2683 #[test]
2684 fn test_write_run_end_encoded() {
2685 let run_ends = Int32Array::from(vec![2, 5, 6]);
2686 let values = StringArray::from(vec![Some("a"), Some("b"), None]);
2687 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2688
2689 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2690 "c1",
2691 ree.data_type().clone(),
2692 true,
2693 )]));
2694
2695 let batch = RecordBatch::try_new(schema, vec![Arc::new(ree)]).unwrap();
2696
2697 let mut buf = Vec::new();
2698 {
2699 let mut writer = LineDelimitedWriter::new(&mut buf);
2700 writer.write_batches(&[&batch]).unwrap();
2701 }
2702
2703 assert_json_eq(
2704 &buf,
2705 r#"{"c1":"a"}
2706{"c1":"a"}
2707{"c1":"b"}
2708{"c1":"b"}
2709{"c1":"b"}
2710{}
2711"#,
2712 );
2713 }
2714
2715 #[test]
2716 fn test_write_run_end_encoded_int_values() {
2717 let run_ends = Int32Array::from(vec![3, 5]);
2718 let values = Int32Array::from(vec![10, 20]);
2719 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2720
2721 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2722 "n",
2723 ree.data_type().clone(),
2724 true,
2725 )]));
2726
2727 let batch = RecordBatch::try_new(schema, vec![Arc::new(ree)]).unwrap();
2728
2729 let json_value: Value = {
2730 let mut buf = Vec::new();
2731 let mut writer = WriterBuilder::new().build::<_, JsonArray>(&mut buf);
2732 writer.write_batches(&[&batch]).unwrap();
2733 writer.finish().unwrap();
2734 serde_json::from_slice(&buf).unwrap()
2735 };
2736
2737 let expected = json!([
2738 {"n": 10},
2739 {"n": 10},
2740 {"n": 10},
2741 {"n": 20},
2742 {"n": 20},
2743 ]);
2744
2745 assert_eq!(json_value, expected);
2746 }
2747
2748 #[test]
2749 fn test_run_end_encoded_roundtrip() {
2750 let run_ends = Int32Array::from(vec![3, 5, 7]);
2751 let values = StringArray::from(vec![Some("a"), None, Some("b")]);
2752 let ree = RunArray::<Int32Type>::try_new(&run_ends, &values).unwrap();
2753
2754 let schema = Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new(
2755 "c",
2756 ree.data_type().clone(),
2757 true,
2758 )]));
2759 let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(ree)]).unwrap();
2760
2761 let mut buf = Vec::new();
2762 {
2763 let mut writer = super::LineDelimitedWriter::new(&mut buf);
2764 writer.write_batches(&[&batch]).unwrap();
2765 }
2766
2767 let batches: Vec<RecordBatch> = ReaderBuilder::new(schema)
2768 .with_batch_size(1024)
2769 .build(std::io::Cursor::new(&buf))
2770 .unwrap()
2771 .collect::<Result<Vec<_>, _>>()
2772 .unwrap();
2773 assert_eq!(batches.len(), 1);
2774
2775 let col = batches[0].column(0);
2776 let run_array = col.as_run::<Int32Type>();
2777
2778 assert_eq!(run_array.len(), 7);
2779 assert_eq!(run_array.run_ends().values(), &[3, 5, 7]);
2780
2781 let values = run_array.values().as_string::<i32>();
2782 assert_eq!(values.len(), 3);
2783 assert_eq!(values.value(0), "a");
2784 assert!(values.is_null(1));
2785 assert_eq!(values.value(2), "b");
2786 }
2787}