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