1use std::borrow::Cow;
137use std::io::BufRead;
138use std::sync::Arc;
139
140use arrow_array::cast::AsArray;
141use arrow_array::timezone::Tz;
142use arrow_array::types::*;
143use arrow_array::{ArrayRef, RecordBatch, RecordBatchReader, downcast_integer};
144use arrow_schema::{ArrowError, DataType, FieldRef, Schema, SchemaRef, TimeUnit};
145use chrono::Utc;
146use serde_core::Serialize;
147
148use crate::StructMode;
149use crate::reader::binary_array::{
150 BinaryArrayDecoder, BinaryViewDecoder, FixedSizeBinaryArrayDecoder,
151};
152use crate::reader::boolean_array::BooleanArrayDecoder;
153use crate::reader::decimal_array::DecimalArrayDecoder;
154use crate::reader::list_array::{
155 FixedSizeListArrayDecoder, ListArrayDecoder, ListViewArrayDecoder,
156};
157use crate::reader::map_array::MapArrayDecoder;
158use crate::reader::null_array::NullArrayDecoder;
159use crate::reader::primitive_array::PrimitiveArrayDecoder;
160use crate::reader::run_end_array::RunEndEncodedArrayDecoder;
161use crate::reader::string_array::StringArrayDecoder;
162use crate::reader::string_view_array::StringViewArrayDecoder;
163use crate::reader::struct_array::StructArrayDecoder;
164use crate::reader::tape::{Tape, TapeDecoder};
165use crate::reader::timestamp_array::TimestampArrayDecoder;
166
167pub use schema::*;
168pub use value_iter::ValueIter;
169
170mod binary_array;
171mod boolean_array;
172mod decimal_array;
173mod list_array;
174mod map_array;
175mod null_array;
176mod primitive_array;
177mod run_end_array;
178mod schema;
179mod serializer;
180mod string_array;
181mod string_view_array;
182mod struct_array;
183mod tape;
184mod timestamp_array;
185mod value_iter;
186
187pub struct ReaderBuilder {
189 batch_size: usize,
190 coerce_primitive: bool,
191 strict_mode: bool,
192 ignore_type_conflicts: bool,
193 is_field: bool,
194 struct_mode: StructMode,
195
196 schema: SchemaRef,
197}
198
199impl ReaderBuilder {
200 pub fn new(schema: SchemaRef) -> Self {
209 Self {
210 batch_size: 1024,
211 coerce_primitive: false,
212 strict_mode: false,
213 ignore_type_conflicts: false,
214 is_field: false,
215 struct_mode: Default::default(),
216 schema,
217 }
218 }
219
220 pub fn new_with_field(field: impl Into<FieldRef>) -> Self {
251 Self {
252 batch_size: 1024,
253 coerce_primitive: false,
254 strict_mode: false,
255 ignore_type_conflicts: false,
256 is_field: true,
257 struct_mode: Default::default(),
258 schema: Arc::new(Schema::new([field.into()])),
259 }
260 }
261
262 pub fn with_batch_size(self, batch_size: usize) -> Self {
264 Self { batch_size, ..self }
265 }
266
267 pub fn with_coerce_primitive(self, coerce_primitive: bool) -> Self {
270 Self {
271 coerce_primitive,
272 ..self
273 }
274 }
275
276 pub fn with_strict_mode(self, strict_mode: bool) -> Self {
282 Self {
283 strict_mode,
284 ..self
285 }
286 }
287
288 pub fn with_struct_mode(self, struct_mode: StructMode) -> Self {
292 Self {
293 struct_mode,
294 ..self
295 }
296 }
297
298 pub fn with_ignore_type_conflicts(self, ignore_type_conflicts: bool) -> Self {
311 Self {
312 ignore_type_conflicts,
313 ..self
314 }
315 }
316
317 pub fn build<R: BufRead>(self, reader: R) -> Result<Reader<R>, ArrowError> {
319 Ok(Reader {
320 reader,
321 decoder: self.build_decoder()?,
322 })
323 }
324
325 pub fn build_decoder(self) -> Result<Decoder, ArrowError> {
327 let (data_type, nullable) = if self.is_field {
328 let field = &self.schema.fields[0];
329 let data_type = Cow::Borrowed(field.data_type());
330 (data_type, field.is_nullable())
331 } else {
332 let data_type = Cow::Owned(DataType::Struct(self.schema.fields.clone()));
333 (data_type, false)
334 };
335
336 let ctx = DecoderContext {
337 coerce_primitive: self.coerce_primitive,
338 strict_mode: self.strict_mode,
339 struct_mode: self.struct_mode,
340 ignore_type_conflicts: self.ignore_type_conflicts,
341 };
342 let decoder = ctx.make_decoder(data_type.as_ref(), nullable)?;
343
344 let num_fields = self.schema.flattened_fields().len();
345
346 Ok(Decoder {
347 decoder,
348 is_field: self.is_field,
349 tape_decoder: TapeDecoder::new(self.batch_size, num_fields),
350 batch_size: self.batch_size,
351 schema: self.schema,
352 })
353 }
354}
355
356pub struct Reader<R> {
360 reader: R,
361 decoder: Decoder,
362}
363
364impl<R> std::fmt::Debug for Reader<R> {
365 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
366 f.debug_struct("Reader")
367 .field("decoder", &self.decoder)
368 .finish()
369 }
370}
371
372impl<R: BufRead> Reader<R> {
373 fn read(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
375 loop {
376 let buf = self.reader.fill_buf()?;
377 if buf.is_empty() {
378 break;
379 }
380 let read = buf.len();
381
382 let decoded = self.decoder.decode(buf)?;
383 self.reader.consume(decoded);
384 if decoded != read {
385 break;
386 }
387 }
388 self.decoder.flush()
389 }
390}
391
392impl<R: BufRead> Iterator for Reader<R> {
393 type Item = Result<RecordBatch, ArrowError>;
394
395 fn next(&mut self) -> Option<Self::Item> {
396 self.read().transpose()
397 }
398}
399
400impl<R: BufRead> RecordBatchReader for Reader<R> {
401 fn schema(&self) -> SchemaRef {
402 self.decoder.schema.clone()
403 }
404}
405
406pub struct Decoder {
447 tape_decoder: TapeDecoder,
448 decoder: Box<dyn ArrayDecoder>,
449 batch_size: usize,
450 is_field: bool,
451 schema: SchemaRef,
452}
453
454impl std::fmt::Debug for Decoder {
455 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
456 f.debug_struct("Decoder")
457 .field("schema", &self.schema)
458 .field("batch_size", &self.batch_size)
459 .finish()
460 }
461}
462
463impl Decoder {
464 pub fn decode(&mut self, buf: &[u8]) -> Result<usize, ArrowError> {
473 self.tape_decoder.decode(buf)
474 }
475
476 pub fn serialize<S: Serialize>(&mut self, rows: &[S]) -> Result<(), ArrowError> {
653 self.tape_decoder.serialize(rows)
654 }
655
656 pub fn has_partial_record(&self) -> bool {
658 self.tape_decoder.has_partial_row()
659 }
660
661 pub fn len(&self) -> usize {
663 self.tape_decoder.num_buffered_rows()
664 }
665
666 pub fn is_empty(&self) -> bool {
668 self.len() == 0
669 }
670
671 pub fn flush(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
678 let tape = self.tape_decoder.finish()?;
679
680 if tape.num_rows() == 0 {
681 return Ok(None);
682 }
683
684 let mut next_object = 1;
686 let pos: Vec<_> = (0..tape.num_rows())
687 .map(|_| {
688 let next = tape.next(next_object, "row").unwrap();
689 std::mem::replace(&mut next_object, next)
690 })
691 .collect();
692
693 let decoded = self.decoder.decode(&tape, &pos)?;
694 self.tape_decoder.clear();
695
696 let batch = match self.is_field {
697 true => RecordBatch::try_new(self.schema.clone(), vec![decoded])?,
698 false => {
699 RecordBatch::from(decoded.as_struct().clone()).with_schema(self.schema.clone())?
700 }
701 };
702
703 Ok(Some(batch))
704 }
705}
706
707trait ArrayDecoder: Send {
708 fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError>;
710}
711
712pub struct DecoderContext {
717 coerce_primitive: bool,
719 strict_mode: bool,
721 struct_mode: StructMode,
723 ignore_type_conflicts: bool,
725}
726
727impl DecoderContext {
728 pub fn coerce_primitive(&self) -> bool {
730 self.coerce_primitive
731 }
732
733 pub fn strict_mode(&self) -> bool {
735 self.strict_mode
736 }
737
738 pub fn struct_mode(&self) -> StructMode {
740 self.struct_mode
741 }
742
743 pub fn ignore_type_conflicts(&self) -> bool {
745 self.ignore_type_conflicts
746 }
747
748 fn make_decoder(
753 &self,
754 data_type: &DataType,
755 is_nullable: bool,
756 ) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
757 make_decoder(self, data_type, is_nullable)
758 }
759}
760
761fn make_decoder(
762 ctx: &DecoderContext,
763 data_type: &DataType,
764 is_nullable: bool,
765) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
766 macro_rules! primitive_decoder {
767 ($t:ty, $data_type:expr) => {
768 Ok(Box::new(PrimitiveArrayDecoder::<$t>::new(ctx, $data_type)))
769 };
770 }
771 macro_rules! timestamp_decoder {
772 ($t:ty, $data_type:expr, $tz:expr) => {{
773 Ok(Box::new(TimestampArrayDecoder::<$t, _>::new(
774 ctx, $data_type, $tz,
775 )))
776 }};
777 }
778 macro_rules! decimal_decoder {
779 ($t:ty, $p:expr, $s:expr) => {
780 Ok(Box::new(DecimalArrayDecoder::<$t>::new(ctx, $p, $s)))
781 };
782 }
783
784 downcast_integer! {
785 *data_type => (primitive_decoder, data_type),
786 DataType::Null => Ok(Box::new(NullArrayDecoder::new(ctx))),
787 DataType::Float16 => primitive_decoder!(Float16Type, data_type),
788 DataType::Float32 => primitive_decoder!(Float32Type, data_type),
789 DataType::Float64 => primitive_decoder!(Float64Type, data_type),
790 DataType::Timestamp(TimeUnit::Second, None) => {
791 timestamp_decoder!(TimestampSecondType, data_type, Utc)
792 },
793 DataType::Timestamp(TimeUnit::Millisecond, None) => {
794 timestamp_decoder!(TimestampMillisecondType, data_type, Utc)
795 },
796 DataType::Timestamp(TimeUnit::Microsecond, None) => {
797 timestamp_decoder!(TimestampMicrosecondType, data_type, Utc)
798 },
799 DataType::Timestamp(TimeUnit::Nanosecond, None) => {
800 timestamp_decoder!(TimestampNanosecondType, data_type, Utc)
801 },
802 DataType::Timestamp(TimeUnit::Second, Some(ref tz)) => {
803 let tz: Tz = tz.parse()?;
804 timestamp_decoder!(TimestampSecondType, data_type, tz)
805 },
806 DataType::Timestamp(TimeUnit::Millisecond, Some(ref tz)) => {
807 let tz: Tz = tz.parse()?;
808 timestamp_decoder!(TimestampMillisecondType, data_type, tz)
809 },
810 DataType::Timestamp(TimeUnit::Microsecond, Some(ref tz)) => {
811 let tz: Tz = tz.parse()?;
812 timestamp_decoder!(TimestampMicrosecondType, data_type, tz)
813 },
814 DataType::Timestamp(TimeUnit::Nanosecond, Some(ref tz)) => {
815 let tz: Tz = tz.parse()?;
816 timestamp_decoder!(TimestampNanosecondType, data_type, tz)
817 },
818 DataType::Date32 => primitive_decoder!(Date32Type, data_type),
819 DataType::Date64 => primitive_decoder!(Date64Type, data_type),
820 DataType::Time32(TimeUnit::Second) => primitive_decoder!(Time32SecondType, data_type),
821 DataType::Time32(TimeUnit::Millisecond) => primitive_decoder!(Time32MillisecondType, data_type),
822 DataType::Time64(TimeUnit::Microsecond) => primitive_decoder!(Time64MicrosecondType, data_type),
823 DataType::Time64(TimeUnit::Nanosecond) => primitive_decoder!(Time64NanosecondType, data_type),
824 DataType::Duration(TimeUnit::Nanosecond) => primitive_decoder!(DurationNanosecondType, data_type),
825 DataType::Duration(TimeUnit::Microsecond) => primitive_decoder!(DurationMicrosecondType, data_type),
826 DataType::Duration(TimeUnit::Millisecond) => primitive_decoder!(DurationMillisecondType, data_type),
827 DataType::Duration(TimeUnit::Second) => primitive_decoder!(DurationSecondType, data_type),
828 DataType::Decimal32(p, s) => decimal_decoder!(Decimal32Type, p, s),
829 DataType::Decimal64(p, s) => decimal_decoder!(Decimal64Type, p, s),
830 DataType::Decimal128(p, s) => decimal_decoder!(Decimal128Type, p, s),
831 DataType::Decimal256(p, s) => decimal_decoder!(Decimal256Type, p, s),
832 DataType::Boolean => Ok(Box::new(BooleanArrayDecoder::new(ctx))),
833 DataType::Utf8 => Ok(Box::new(StringArrayDecoder::<i32>::new(ctx))),
834 DataType::Utf8View => Ok(Box::new(StringViewArrayDecoder::new(ctx))),
835 DataType::LargeUtf8 => Ok(Box::new(StringArrayDecoder::<i64>::new(ctx))),
836 DataType::List(_) => Ok(Box::new(ListArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
837 DataType::LargeList(_) => Ok(Box::new(ListArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
838 DataType::ListView(_) => Ok(Box::new(ListViewArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
839 DataType::LargeListView(_) => Ok(Box::new(ListViewArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
840 DataType::FixedSizeList(_, _) => Ok(Box::new(FixedSizeListArrayDecoder::new(ctx, data_type, is_nullable)?)),
841 DataType::Struct(_) => Ok(Box::new(StructArrayDecoder::new(ctx, data_type, is_nullable)?)),
842 DataType::Binary => Ok(Box::new(BinaryArrayDecoder::<i32>::default())),
843 DataType::LargeBinary => Ok(Box::new(BinaryArrayDecoder::<i64>::default())),
844 DataType::FixedSizeBinary(len) => Ok(Box::new(FixedSizeBinaryArrayDecoder::new(len))),
845 DataType::BinaryView => Ok(Box::new(BinaryViewDecoder::default())),
846 DataType::Map(_, _) => Ok(Box::new(MapArrayDecoder::new(ctx, data_type, is_nullable)?)),
847 DataType::RunEndEncoded(ref r, _) => match r.data_type() {
848 DataType::Int16 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int16Type>::new(ctx, data_type, is_nullable)?)),
849 DataType::Int32 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int32Type>::new(ctx, data_type, is_nullable)?)),
850 DataType::Int64 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int64Type>::new(ctx, data_type, is_nullable)?)),
851 d => unreachable!("unsupported run end index type: {d}"),
852 },
853 _ => Err(ArrowError::NotYetImplemented(format!("Support for {data_type} in JSON reader")))
854 }
855}
856
857#[cfg(test)]
858mod tests {
859 use arrow_array::cast::AsArray;
860 use arrow_array::{
861 Array, BooleanArray, Float64Array, GenericListViewArray, Int32Array, ListArray, MapArray,
862 NullArray, OffsetSizeTrait, StringArray, StringViewArray, StructArray,
863 };
864 use arrow_buffer::{ArrowNativeType, NullBuffer, OffsetBuffer, ScalarBuffer};
865 use arrow_cast::display::{ArrayFormatter, FormatOptions};
866 use arrow_schema::{Field, Fields};
867 use serde_json::json;
868 use std::fs::File;
869 use std::io::{BufReader, Cursor, Seek};
870
871 use super::*;
872
873 fn do_read(
874 buf: &str,
875 batch_size: usize,
876 coerce_primitive: bool,
877 strict_mode: bool,
878 schema: SchemaRef,
879 ) -> Vec<RecordBatch> {
880 let mut unbuffered = vec![];
881
882 for batch_size in [1, 3, 100, batch_size] {
884 unbuffered = ReaderBuilder::new(schema.clone())
885 .with_batch_size(batch_size)
886 .with_coerce_primitive(coerce_primitive)
887 .build(Cursor::new(buf.as_bytes()))
888 .unwrap()
889 .collect::<Result<Vec<_>, _>>()
890 .unwrap();
891
892 for b in unbuffered.iter().take(unbuffered.len() - 1) {
893 assert_eq!(b.num_rows(), batch_size)
894 }
895
896 for b in [1, 3, 5] {
898 let buffered = ReaderBuilder::new(schema.clone())
899 .with_batch_size(batch_size)
900 .with_coerce_primitive(coerce_primitive)
901 .with_strict_mode(strict_mode)
902 .build(BufReader::with_capacity(b, Cursor::new(buf.as_bytes())))
903 .unwrap()
904 .collect::<Result<Vec<_>, _>>()
905 .unwrap();
906 assert_eq!(unbuffered, buffered);
907 }
908 }
909
910 unbuffered
911 }
912
913 #[test]
914 fn test_basic() {
915 let buf = r#"
916 {"a": 1, "b": 2, "c": true, "d": 1}
917 {"a": 2E0, "b": 4, "c": false, "d": 2, "e": 254}
918
919 {"b": 6, "a": 2.0, "d": 45}
920 {"b": "5", "a": 2}
921 {"b": 4e0}
922 {"b": 7, "a": null}
923 "#;
924
925 let schema = Arc::new(Schema::new(vec![
926 Field::new("a", DataType::Int64, true),
927 Field::new("b", DataType::Int32, true),
928 Field::new("c", DataType::Boolean, true),
929 Field::new("d", DataType::Date32, true),
930 Field::new("e", DataType::Date64, true),
931 ]));
932
933 let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
934 assert!(decoder.is_empty());
935 assert_eq!(decoder.len(), 0);
936 assert!(!decoder.has_partial_record());
937 assert_eq!(decoder.decode(buf.as_bytes()).unwrap(), 221);
938 assert!(!decoder.is_empty());
939 assert_eq!(decoder.len(), 6);
940 assert!(!decoder.has_partial_record());
941 let batch = decoder.flush().unwrap().unwrap();
942 assert_eq!(batch.num_rows(), 6);
943 assert!(decoder.is_empty());
944 assert_eq!(decoder.len(), 0);
945 assert!(!decoder.has_partial_record());
946
947 let batches = do_read(buf, 1024, false, false, schema);
948 assert_eq!(batches.len(), 1);
949
950 let col1 = batches[0].column(0).as_primitive::<Int64Type>();
951 assert_eq!(col1.null_count(), 2);
952 assert_eq!(col1.values(), &[1, 2, 2, 2, 0, 0]);
953 assert!(col1.is_null(4));
954 assert!(col1.is_null(5));
955
956 let col2 = batches[0].column(1).as_primitive::<Int32Type>();
957 assert_eq!(col2.null_count(), 0);
958 assert_eq!(col2.values(), &[2, 4, 6, 5, 4, 7]);
959
960 let col3 = batches[0].column(2).as_boolean();
961 assert_eq!(col3.null_count(), 4);
962 assert!(col3.value(0));
963 assert!(!col3.is_null(0));
964 assert!(!col3.value(1));
965 assert!(!col3.is_null(1));
966
967 let col4 = batches[0].column(3).as_primitive::<Date32Type>();
968 assert_eq!(col4.null_count(), 3);
969 assert!(col4.is_null(3));
970 assert_eq!(col4.values(), &[1, 2, 45, 0, 0, 0]);
971
972 let col5 = batches[0].column(4).as_primitive::<Date64Type>();
973 assert_eq!(col5.null_count(), 5);
974 assert!(col5.is_null(0));
975 assert!(col5.is_null(2));
976 assert!(col5.is_null(3));
977 assert_eq!(col5.values(), &[0, 254, 0, 0, 0, 0]);
978 }
979
980 #[test]
981 fn test_string() {
982 let buf = r#"
983 {"a": "1", "b": "2"}
984 {"a": "hello", "b": "shoo"}
985 {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
986
987 {"b": null}
988 {"b": "", "a": null}
989
990 "#;
991 let schema = Arc::new(Schema::new(vec![
992 Field::new("a", DataType::Utf8, true),
993 Field::new("b", DataType::LargeUtf8, true),
994 ]));
995
996 let batches = do_read(buf, 1024, false, false, schema);
997 assert_eq!(batches.len(), 1);
998
999 let col1 = batches[0].column(0).as_string::<i32>();
1000 assert_eq!(col1.null_count(), 2);
1001 assert_eq!(col1.value(0), "1");
1002 assert_eq!(col1.value(1), "hello");
1003 assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1004 assert!(col1.is_null(3));
1005 assert!(col1.is_null(4));
1006
1007 let col2 = batches[0].column(1).as_string::<i64>();
1008 assert_eq!(col2.null_count(), 1);
1009 assert_eq!(col2.value(0), "2");
1010 assert_eq!(col2.value(1), "shoo");
1011 assert_eq!(col2.value(2), "\t😁foo");
1012 assert!(col2.is_null(3));
1013 assert_eq!(col2.value(4), "");
1014 }
1015
1016 #[test]
1017 fn test_long_string_view_allocation() {
1018 let expected_capacity: usize = 41;
1028
1029 let buf = r#"
1030 {"a": "short", "b": "dummy"}
1031 {"a": "this is definitely long", "b": "dummy"}
1032 {"a": "hello", "b": "dummy"}
1033 {"a": "\nfoobar😀asfgÿ", "b": "dummy"}
1034 "#;
1035
1036 let schema = Arc::new(Schema::new(vec![
1037 Field::new("a", DataType::Utf8View, true),
1038 Field::new("b", DataType::LargeUtf8, true),
1039 ]));
1040
1041 let batches = do_read(buf, 1024, false, false, schema);
1042 assert_eq!(batches.len(), 1, "Expected one record batch");
1043
1044 let col_a = batches[0].column(0);
1046 let string_view_array = col_a
1047 .as_any()
1048 .downcast_ref::<StringViewArray>()
1049 .expect("Column should be a StringViewArray");
1050
1051 let data_buffer = string_view_array.to_data().buffers()[0].len();
1054
1055 assert!(
1058 data_buffer >= expected_capacity,
1059 "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1060 );
1061
1062 assert_eq!(string_view_array.value(0), "short");
1064 assert_eq!(string_view_array.value(1), "this is definitely long");
1065 assert_eq!(string_view_array.value(2), "hello");
1066 assert_eq!(string_view_array.value(3), "\nfoobar😀asfgÿ");
1067 }
1068
1069 #[test]
1071 fn test_numeric_view_allocation() {
1072 let expected_capacity: usize = 33;
1080
1081 let buf = r#"
1082 {"n": 123456789}
1083 {"n": 1000000000000}
1084 {"n": 3.1415}
1085 {"n": 2.718281828459045}
1086 "#;
1087
1088 let schema = Arc::new(Schema::new(vec![Field::new("n", DataType::Utf8View, true)]));
1089
1090 let batches = do_read(buf, 1024, true, false, schema);
1091 assert_eq!(batches.len(), 1, "Expected one record batch");
1092
1093 let col_n = batches[0].column(0);
1094 let string_view_array = col_n
1095 .as_any()
1096 .downcast_ref::<StringViewArray>()
1097 .expect("Column should be a StringViewArray");
1098
1099 let data_buffer = string_view_array.to_data().buffers()[0].len();
1101 assert!(
1102 data_buffer >= expected_capacity,
1103 "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1104 );
1105
1106 assert_eq!(string_view_array.value(0), "123456789");
1109 assert_eq!(string_view_array.value(1), "1000000000000");
1110 assert_eq!(string_view_array.value(2), "3.1415");
1111 assert_eq!(string_view_array.value(3), "2.718281828459045");
1112 }
1113
1114 #[test]
1115 fn test_string_with_uft8view() {
1116 let buf = r#"
1117 {"a": "1", "b": "2"}
1118 {"a": "hello", "b": "shoo"}
1119 {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
1120
1121 {"b": null}
1122 {"b": "", "a": null}
1123
1124 "#;
1125 let schema = Arc::new(Schema::new(vec![
1126 Field::new("a", DataType::Utf8View, true),
1127 Field::new("b", DataType::LargeUtf8, true),
1128 ]));
1129
1130 let batches = do_read(buf, 1024, false, false, schema);
1131 assert_eq!(batches.len(), 1);
1132
1133 let col1 = batches[0].column(0).as_string_view();
1134 assert_eq!(col1.null_count(), 2);
1135 assert_eq!(col1.value(0), "1");
1136 assert_eq!(col1.value(1), "hello");
1137 assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1138 assert!(col1.is_null(3));
1139 assert!(col1.is_null(4));
1140 assert_eq!(col1.data_type(), &DataType::Utf8View);
1141
1142 let col2 = batches[0].column(1).as_string::<i64>();
1143 assert_eq!(col2.null_count(), 1);
1144 assert_eq!(col2.value(0), "2");
1145 assert_eq!(col2.value(1), "shoo");
1146 assert_eq!(col2.value(2), "\t😁foo");
1147 assert!(col2.is_null(3));
1148 assert_eq!(col2.value(4), "");
1149 }
1150
1151 #[test]
1152 fn test_complex() {
1153 let buf = r#"
1154 {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3}, {"c": 4}]}}
1155 {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1156 {"list": null, "nested": {"a": null}}
1157 "#;
1158
1159 let schema = Arc::new(Schema::new(vec![
1160 Field::new_list("list", Field::new("element", DataType::Int32, false), true),
1161 Field::new_struct(
1162 "nested",
1163 vec![
1164 Field::new("a", DataType::Int32, true),
1165 Field::new("b", DataType::Int32, true),
1166 ],
1167 true,
1168 ),
1169 Field::new_struct(
1170 "nested_list",
1171 vec![Field::new_list(
1172 "list2",
1173 Field::new_struct(
1174 "element",
1175 vec![Field::new("c", DataType::Int32, false)],
1176 false,
1177 ),
1178 true,
1179 )],
1180 true,
1181 ),
1182 ]));
1183
1184 let batches = do_read(buf, 1024, false, false, schema);
1185 assert_eq!(batches.len(), 1);
1186
1187 let list = batches[0].column(0).as_list::<i32>();
1188 assert_eq!(list.len(), 3);
1189 assert_eq!(list.value_offsets(), &[0, 0, 2, 2]);
1190 assert_eq!(list.null_count(), 1);
1191 assert!(list.is_null(2));
1192 let list_values = list.values().as_primitive::<Int32Type>();
1193 assert_eq!(list_values.values(), &[5, 6]);
1194
1195 let nested = batches[0].column(1).as_struct();
1196 let a = nested.column(0).as_primitive::<Int32Type>();
1197 assert_eq!(list.null_count(), 1);
1198 assert_eq!(a.values(), &[1, 7, 0]);
1199 assert!(list.is_null(2));
1200
1201 let b = nested.column(1).as_primitive::<Int32Type>();
1202 assert_eq!(b.null_count(), 2);
1203 assert_eq!(b.len(), 3);
1204 assert_eq!(b.value(0), 2);
1205 assert!(b.is_null(1));
1206 assert!(b.is_null(2));
1207
1208 let nested_list = batches[0].column(2).as_struct();
1209 assert_eq!(nested_list.len(), 3);
1210 assert_eq!(nested_list.null_count(), 1);
1211 assert!(nested_list.is_null(2));
1212
1213 let list2 = nested_list.column(0).as_list::<i32>();
1214 assert_eq!(list2.len(), 3);
1215 assert_eq!(list2.null_count(), 1);
1216 assert_eq!(list2.value_offsets(), &[0, 2, 2, 2]);
1217 assert!(list2.is_null(2));
1218
1219 let list2_values = list2.values().as_struct();
1220
1221 let c = list2_values.column(0).as_primitive::<Int32Type>();
1222 assert_eq!(c.values(), &[3, 4]);
1223 }
1224
1225 #[test]
1226 fn test_projection() {
1227 let buf = r#"
1228 {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3, "d": 5}, {"c": 4}]}}
1229 {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1230 "#;
1231
1232 let schema = Arc::new(Schema::new(vec![
1233 Field::new_struct(
1234 "nested",
1235 vec![Field::new("a", DataType::Int32, false)],
1236 true,
1237 ),
1238 Field::new_struct(
1239 "nested_list",
1240 vec![Field::new_list(
1241 "list2",
1242 Field::new_struct(
1243 "element",
1244 vec![Field::new("d", DataType::Int32, true)],
1245 false,
1246 ),
1247 true,
1248 )],
1249 true,
1250 ),
1251 ]));
1252
1253 let batches = do_read(buf, 1024, false, false, schema);
1254 assert_eq!(batches.len(), 1);
1255
1256 let nested = batches[0].column(0).as_struct();
1257 assert_eq!(nested.num_columns(), 1);
1258 let a = nested.column(0).as_primitive::<Int32Type>();
1259 assert_eq!(a.null_count(), 0);
1260 assert_eq!(a.values(), &[1, 7]);
1261
1262 let nested_list = batches[0].column(1).as_struct();
1263 assert_eq!(nested_list.num_columns(), 1);
1264 assert_eq!(nested_list.null_count(), 0);
1265
1266 let list2 = nested_list.column(0).as_list::<i32>();
1267 assert_eq!(list2.value_offsets(), &[0, 2, 2]);
1268 assert_eq!(list2.null_count(), 0);
1269
1270 let child = list2.values().as_struct();
1271 assert_eq!(child.num_columns(), 1);
1272 assert_eq!(child.len(), 2);
1273 assert_eq!(child.null_count(), 0);
1274
1275 let c = child.column(0).as_primitive::<Int32Type>();
1276 assert_eq!(c.values(), &[5, 0]);
1277 assert_eq!(c.null_count(), 1);
1278 assert!(c.is_null(1));
1279 }
1280
1281 #[test]
1282 fn test_map() {
1283 let buf = r#"
1284 {"map": {"a": ["foo", null]}}
1285 {"map": {"a": [null], "b": []}}
1286 {"map": {"c": null, "a": ["baz"]}}
1287 "#;
1288 let map = Field::new_map(
1289 "map",
1290 "entries",
1291 Field::new("key", DataType::Utf8, false),
1292 Field::new_list("value", Field::new("element", DataType::Utf8, true), true),
1293 false,
1294 true,
1295 );
1296
1297 let schema = Arc::new(Schema::new(vec![map]));
1298
1299 let batches = do_read(buf, 1024, false, false, schema);
1300 assert_eq!(batches.len(), 1);
1301
1302 let map = batches[0].column(0).as_map();
1303 let map_keys = map.keys().as_string::<i32>();
1304 let map_values = map.values().as_list::<i32>();
1305 assert_eq!(map.value_offsets(), &[0, 1, 3, 5]);
1306
1307 let k: Vec<_> = map_keys.iter().flatten().collect();
1308 assert_eq!(&k, &["a", "a", "b", "c", "a"]);
1309
1310 let list_values = map_values.values().as_string::<i32>();
1311 let lv: Vec<_> = list_values.iter().collect();
1312 assert_eq!(&lv, &[Some("foo"), None, None, Some("baz")]);
1313 assert_eq!(map_values.value_offsets(), &[0, 2, 3, 3, 3, 4]);
1314 assert_eq!(map_values.null_count(), 1);
1315 assert!(map_values.is_null(3));
1316
1317 let options = FormatOptions::default().with_null("null");
1318 let formatter = ArrayFormatter::try_new(map, &options).unwrap();
1319 assert_eq!(formatter.value(0).to_string(), "{a: [foo, null]}");
1320 assert_eq!(formatter.value(1).to_string(), "{a: [null], b: []}");
1321 assert_eq!(formatter.value(2).to_string(), "{c: null, a: [baz]}");
1322 }
1323
1324 #[test]
1325 fn test_map_non_nullable_value() {
1326 let map = Field::new_map(
1327 "map",
1328 "entries",
1329 Field::new("keys", DataType::Utf8, false),
1330 Field::new("values", DataType::Utf8, false),
1331 false,
1332 false,
1333 );
1334 let schema = Arc::new(Schema::new(vec![map]));
1335 let buf = r#"{"map": {"key": null}}"#;
1336
1337 let err = ReaderBuilder::new(schema)
1338 .build(Cursor::new(buf.as_bytes()))
1339 .unwrap()
1340 .read()
1341 .unwrap_err();
1342
1343 assert_eq!(
1344 err.to_string(),
1345 "Invalid argument error: Found unmasked nulls for non-nullable StructArray field \"values\""
1346 );
1347 }
1348
1349 #[test]
1350 fn test_not_coercing_primitive_into_string_without_flag() {
1351 let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, true)]));
1352
1353 let buf = r#"{"a": 1}"#;
1354 let err = ReaderBuilder::new(schema.clone())
1355 .with_batch_size(1024)
1356 .build(Cursor::new(buf.as_bytes()))
1357 .unwrap()
1358 .read()
1359 .unwrap_err();
1360
1361 assert_eq!(
1362 err.to_string(),
1363 "Json error: whilst decoding field 'a': expected string got 1"
1364 );
1365
1366 let buf = r#"{"a": true}"#;
1367 let err = ReaderBuilder::new(schema)
1368 .with_batch_size(1024)
1369 .build(Cursor::new(buf.as_bytes()))
1370 .unwrap()
1371 .read()
1372 .unwrap_err();
1373
1374 assert_eq!(
1375 err.to_string(),
1376 "Json error: whilst decoding field 'a': expected string got true"
1377 );
1378 }
1379
1380 #[test]
1381 fn test_coercing_primitive_into_string() {
1382 let buf = r#"
1383 {"a": 1, "b": 2, "c": true}
1384 {"a": 2E0, "b": 4, "c": false}
1385
1386 {"b": 6, "a": 2.0}
1387 {"b": "5", "a": 2}
1388 {"b": 4e0}
1389 {"b": 7, "a": null}
1390 "#;
1391
1392 let schema = Arc::new(Schema::new(vec![
1393 Field::new("a", DataType::Utf8, true),
1394 Field::new("b", DataType::Utf8, true),
1395 Field::new("c", DataType::Utf8, true),
1396 ]));
1397
1398 let batches = do_read(buf, 1024, true, false, schema);
1399 assert_eq!(batches.len(), 1);
1400
1401 let col1 = batches[0].column(0).as_string::<i32>();
1402 assert_eq!(col1.null_count(), 2);
1403 assert_eq!(col1.value(0), "1");
1404 assert_eq!(col1.value(1), "2E0");
1405 assert_eq!(col1.value(2), "2.0");
1406 assert_eq!(col1.value(3), "2");
1407 assert!(col1.is_null(4));
1408 assert!(col1.is_null(5));
1409
1410 let col2 = batches[0].column(1).as_string::<i32>();
1411 assert_eq!(col2.null_count(), 0);
1412 assert_eq!(col2.value(0), "2");
1413 assert_eq!(col2.value(1), "4");
1414 assert_eq!(col2.value(2), "6");
1415 assert_eq!(col2.value(3), "5");
1416 assert_eq!(col2.value(4), "4e0");
1417 assert_eq!(col2.value(5), "7");
1418
1419 let col3 = batches[0].column(2).as_string::<i32>();
1420 assert_eq!(col3.null_count(), 4);
1421 assert_eq!(col3.value(0), "true");
1422 assert_eq!(col3.value(1), "false");
1423 assert!(col3.is_null(2));
1424 assert!(col3.is_null(3));
1425 assert!(col3.is_null(4));
1426 assert!(col3.is_null(5));
1427 }
1428
1429 fn test_decimal<T: DecimalType>(data_type: DataType) {
1430 let buf = r#"
1431 {"a": 1, "b": 2, "c": 38.30}
1432 {"a": 2, "b": 4, "c": 123.456}
1433
1434 {"b": 1337, "a": "2.0452"}
1435 {"b": "5", "a": "11034.2"}
1436 {"b": 40}
1437 {"b": 1234, "a": null}
1438 "#;
1439
1440 let schema = Arc::new(Schema::new(vec![
1441 Field::new("a", data_type.clone(), true),
1442 Field::new("b", data_type.clone(), true),
1443 Field::new("c", data_type, true),
1444 ]));
1445
1446 let batches = do_read(buf, 1024, true, false, schema);
1447 assert_eq!(batches.len(), 1);
1448
1449 let col1 = batches[0].column(0).as_primitive::<T>();
1450 assert_eq!(col1.null_count(), 2);
1451 assert!(col1.is_null(4));
1452 assert!(col1.is_null(5));
1453 assert_eq!(
1454 col1.values(),
1455 &[100, 200, 204, 1103420, 0, 0].map(T::Native::usize_as)
1456 );
1457
1458 let col2 = batches[0].column(1).as_primitive::<T>();
1459 assert_eq!(col2.null_count(), 0);
1460 assert_eq!(
1461 col2.values(),
1462 &[200, 400, 133700, 500, 4000, 123400].map(T::Native::usize_as)
1463 );
1464
1465 let col3 = batches[0].column(2).as_primitive::<T>();
1466 assert_eq!(col3.null_count(), 4);
1467 assert!(!col3.is_null(0));
1468 assert!(!col3.is_null(1));
1469 assert!(col3.is_null(2));
1470 assert!(col3.is_null(3));
1471 assert!(col3.is_null(4));
1472 assert!(col3.is_null(5));
1473 assert_eq!(
1474 col3.values(),
1475 &[3830, 12345, 0, 0, 0, 0].map(T::Native::usize_as)
1476 );
1477 }
1478
1479 #[test]
1480 fn test_decimals() {
1481 test_decimal::<Decimal32Type>(DataType::Decimal32(8, 2));
1482 test_decimal::<Decimal64Type>(DataType::Decimal64(10, 2));
1483 test_decimal::<Decimal128Type>(DataType::Decimal128(10, 2));
1484 test_decimal::<Decimal256Type>(DataType::Decimal256(10, 2));
1485 }
1486
1487 fn test_timestamp<T: ArrowTimestampType>() {
1488 let buf = r#"
1489 {"a": 1, "b": "2020-09-08T13:42:29.190855+00:00", "c": 38.30, "d": "1997-01-31T09:26:56.123"}
1490 {"a": 2, "b": "2020-09-08T13:42:29.190855Z", "c": 123.456, "d": 123.456}
1491
1492 {"b": 1337, "b": "2020-09-08T13:42:29Z", "c": "1997-01-31T09:26:56.123", "d": "1997-01-31T09:26:56.123Z"}
1493 {"b": 40, "c": "2020-09-08T13:42:29.190855+00:00", "d": "1997-01-31 09:26:56.123-05:00"}
1494 {"b": 1234, "a": null, "c": "1997-01-31 09:26:56.123Z", "d": "1997-01-31 092656"}
1495 {"c": "1997-01-31T14:26:56.123-05:00", "d": "1997-01-31"}
1496 "#;
1497
1498 let with_timezone = DataType::Timestamp(T::UNIT, Some("+08:00".into()));
1499 let schema = Arc::new(Schema::new(vec![
1500 Field::new("a", T::DATA_TYPE, true),
1501 Field::new("b", T::DATA_TYPE, true),
1502 Field::new("c", T::DATA_TYPE, true),
1503 Field::new("d", with_timezone, true),
1504 ]));
1505
1506 let batches = do_read(buf, 1024, true, false, schema);
1507 assert_eq!(batches.len(), 1);
1508
1509 let unit_in_nanos: i64 = match T::UNIT {
1510 TimeUnit::Second => 1_000_000_000,
1511 TimeUnit::Millisecond => 1_000_000,
1512 TimeUnit::Microsecond => 1_000,
1513 TimeUnit::Nanosecond => 1,
1514 };
1515
1516 let col1 = batches[0].column(0).as_primitive::<T>();
1517 assert_eq!(col1.null_count(), 4);
1518 assert!(col1.is_null(2));
1519 assert!(col1.is_null(3));
1520 assert!(col1.is_null(4));
1521 assert!(col1.is_null(5));
1522 assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1523
1524 let col2 = batches[0].column(1).as_primitive::<T>();
1525 assert_eq!(col2.null_count(), 1);
1526 assert!(col2.is_null(5));
1527 assert_eq!(
1528 col2.values(),
1529 &[
1530 1599572549190855000 / unit_in_nanos,
1531 1599572549190855000 / unit_in_nanos,
1532 1599572549000000000 / unit_in_nanos,
1533 40,
1534 1234,
1535 0
1536 ]
1537 );
1538
1539 let col3 = batches[0].column(2).as_primitive::<T>();
1540 assert_eq!(col3.null_count(), 0);
1541 assert_eq!(
1542 col3.values(),
1543 &[
1544 38,
1545 123,
1546 854702816123000000 / unit_in_nanos,
1547 1599572549190855000 / unit_in_nanos,
1548 854702816123000000 / unit_in_nanos,
1549 854738816123000000 / unit_in_nanos
1550 ]
1551 );
1552
1553 let col4 = batches[0].column(3).as_primitive::<T>();
1554
1555 assert_eq!(col4.null_count(), 0);
1556 assert_eq!(
1557 col4.values(),
1558 &[
1559 854674016123000000 / unit_in_nanos,
1560 123,
1561 854702816123000000 / unit_in_nanos,
1562 854720816123000000 / unit_in_nanos,
1563 854674016000000000 / unit_in_nanos,
1564 854640000000000000 / unit_in_nanos
1565 ]
1566 );
1567 }
1568
1569 #[test]
1570 fn test_timestamps() {
1571 test_timestamp::<TimestampSecondType>();
1572 test_timestamp::<TimestampMillisecondType>();
1573 test_timestamp::<TimestampMicrosecondType>();
1574 test_timestamp::<TimestampNanosecondType>();
1575 }
1576
1577 fn test_time<T: ArrowTemporalType>() {
1578 let buf = r#"
1579 {"a": 1, "b": "09:26:56.123 AM", "c": 38.30}
1580 {"a": 2, "b": "23:59:59", "c": 123.456}
1581
1582 {"b": 1337, "b": "6:00 pm", "c": "09:26:56.123"}
1583 {"b": 40, "c": "13:42:29.190855"}
1584 {"b": 1234, "a": null, "c": "09:26:56.123"}
1585 {"c": "14:26:56.123"}
1586 "#;
1587
1588 let unit = match T::DATA_TYPE {
1589 DataType::Time32(unit) | DataType::Time64(unit) => unit,
1590 _ => unreachable!(),
1591 };
1592
1593 let unit_in_nanos = match unit {
1594 TimeUnit::Second => 1_000_000_000,
1595 TimeUnit::Millisecond => 1_000_000,
1596 TimeUnit::Microsecond => 1_000,
1597 TimeUnit::Nanosecond => 1,
1598 };
1599
1600 let schema = Arc::new(Schema::new(vec![
1601 Field::new("a", T::DATA_TYPE, true),
1602 Field::new("b", T::DATA_TYPE, true),
1603 Field::new("c", T::DATA_TYPE, true),
1604 ]));
1605
1606 let batches = do_read(buf, 1024, true, false, schema);
1607 assert_eq!(batches.len(), 1);
1608
1609 let col1 = batches[0].column(0).as_primitive::<T>();
1610 assert_eq!(col1.null_count(), 4);
1611 assert!(col1.is_null(2));
1612 assert!(col1.is_null(3));
1613 assert!(col1.is_null(4));
1614 assert!(col1.is_null(5));
1615 assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1616
1617 let col2 = batches[0].column(1).as_primitive::<T>();
1618 assert_eq!(col2.null_count(), 1);
1619 assert!(col2.is_null(5));
1620 assert_eq!(
1621 col2.values(),
1622 &[
1623 34016123000000 / unit_in_nanos,
1624 86399000000000 / unit_in_nanos,
1625 64800000000000 / unit_in_nanos,
1626 40,
1627 1234,
1628 0
1629 ]
1630 .map(T::Native::usize_as)
1631 );
1632
1633 let col3 = batches[0].column(2).as_primitive::<T>();
1634 assert_eq!(col3.null_count(), 0);
1635 assert_eq!(
1636 col3.values(),
1637 &[
1638 38,
1639 123,
1640 34016123000000 / unit_in_nanos,
1641 49349190855000 / unit_in_nanos,
1642 34016123000000 / unit_in_nanos,
1643 52016123000000 / unit_in_nanos
1644 ]
1645 .map(T::Native::usize_as)
1646 );
1647 }
1648
1649 #[test]
1650 fn test_times() {
1651 test_time::<Time32MillisecondType>();
1652 test_time::<Time32SecondType>();
1653 test_time::<Time64MicrosecondType>();
1654 test_time::<Time64NanosecondType>();
1655 }
1656
1657 fn test_duration<T: ArrowTemporalType>() {
1658 let buf = r#"
1659 {"a": 1, "b": "2"}
1660 {"a": 3, "b": null}
1661 "#;
1662
1663 let schema = Arc::new(Schema::new(vec![
1664 Field::new("a", T::DATA_TYPE, true),
1665 Field::new("b", T::DATA_TYPE, true),
1666 ]));
1667
1668 let batches = do_read(buf, 1024, true, false, schema);
1669 assert_eq!(batches.len(), 1);
1670
1671 let col_a = batches[0].column_by_name("a").unwrap().as_primitive::<T>();
1672 assert_eq!(col_a.null_count(), 0);
1673 assert_eq!(col_a.values(), &[1, 3].map(T::Native::usize_as));
1674
1675 let col2 = batches[0].column_by_name("b").unwrap().as_primitive::<T>();
1676 assert_eq!(col2.null_count(), 1);
1677 assert_eq!(col2.values(), &[2, 0].map(T::Native::usize_as));
1678 }
1679
1680 #[test]
1681 fn test_durations() {
1682 test_duration::<DurationNanosecondType>();
1683 test_duration::<DurationMicrosecondType>();
1684 test_duration::<DurationMillisecondType>();
1685 test_duration::<DurationSecondType>();
1686 }
1687
1688 #[test]
1689 fn test_delta_checkpoint() {
1690 let json = "{\"protocol\":{\"minReaderVersion\":1,\"minWriterVersion\":2}}";
1691 let schema = Arc::new(Schema::new(vec![
1692 Field::new_struct(
1693 "protocol",
1694 vec![
1695 Field::new("minReaderVersion", DataType::Int32, true),
1696 Field::new("minWriterVersion", DataType::Int32, true),
1697 ],
1698 true,
1699 ),
1700 Field::new_struct(
1701 "add",
1702 vec![Field::new_map(
1703 "partitionValues",
1704 "key_value",
1705 Field::new("key", DataType::Utf8, false),
1706 Field::new("value", DataType::Utf8, true),
1707 false,
1708 false,
1709 )],
1710 true,
1711 ),
1712 ]));
1713
1714 let batches = do_read(json, 1024, true, false, schema);
1715 assert_eq!(batches.len(), 1);
1716
1717 let s: StructArray = batches.into_iter().next().unwrap().into();
1718 let opts = FormatOptions::default().with_null("null");
1719 let formatter = ArrayFormatter::try_new(&s, &opts).unwrap();
1720 assert_eq!(
1721 formatter.value(0).to_string(),
1722 "{protocol: {minReaderVersion: 1, minWriterVersion: 2}, add: null}"
1723 );
1724 }
1725
1726 #[test]
1727 fn struct_nullability() {
1728 let do_test = |child: DataType| {
1729 let non_null = r#"{"foo": {}}"#;
1731 let schema = Arc::new(Schema::new(vec![Field::new_struct(
1732 "foo",
1733 vec![Field::new("bar", child, false)],
1734 true,
1735 )]));
1736 let mut reader = ReaderBuilder::new(schema.clone())
1737 .build(Cursor::new(non_null.as_bytes()))
1738 .unwrap();
1739 assert!(reader.next().unwrap().is_err()); let null = r#"{"foo": {bar: null}}"#;
1742 let mut reader = ReaderBuilder::new(schema.clone())
1743 .build(Cursor::new(null.as_bytes()))
1744 .unwrap();
1745 assert!(reader.next().unwrap().is_err()); let null = r#"{"foo": null}"#;
1749 let mut reader = ReaderBuilder::new(schema)
1750 .build(Cursor::new(null.as_bytes()))
1751 .unwrap();
1752 let batch = reader.next().unwrap().unwrap();
1753 assert_eq!(batch.num_columns(), 1);
1754 let foo = batch.column(0).as_struct();
1755 assert_eq!(foo.len(), 1);
1756 assert!(foo.is_null(0));
1757 assert_eq!(foo.num_columns(), 1);
1758
1759 let bar = foo.column(0);
1760 assert_eq!(bar.len(), 1);
1761 assert!(bar.is_null(0));
1763 };
1764
1765 do_test(DataType::Boolean);
1766 do_test(DataType::Int32);
1767 do_test(DataType::Utf8);
1768 do_test(DataType::Decimal128(2, 1));
1769 do_test(DataType::Timestamp(
1770 TimeUnit::Microsecond,
1771 Some("+00:00".into()),
1772 ));
1773 }
1774
1775 #[test]
1776 fn test_truncation() {
1777 let buf = r#"
1778 {"i64": 9223372036854775807, "u64": 18446744073709551615 }
1779 {"i64": "9223372036854775807", "u64": "18446744073709551615" }
1780 {"i64": -9223372036854775808, "u64": 0 }
1781 {"i64": "-9223372036854775808", "u64": 0 }
1782 "#;
1783
1784 let schema = Arc::new(Schema::new(vec![
1785 Field::new("i64", DataType::Int64, true),
1786 Field::new("u64", DataType::UInt64, true),
1787 ]));
1788
1789 let batches = do_read(buf, 1024, true, false, schema);
1790 assert_eq!(batches.len(), 1);
1791
1792 let i64 = batches[0].column(0).as_primitive::<Int64Type>();
1793 assert_eq!(i64.values(), &[i64::MAX, i64::MAX, i64::MIN, i64::MIN]);
1794
1795 let u64 = batches[0].column(1).as_primitive::<UInt64Type>();
1796 assert_eq!(u64.values(), &[u64::MAX, u64::MAX, u64::MIN, u64::MIN]);
1797 }
1798
1799 #[test]
1800 fn test_timestamp_truncation() {
1801 let buf = r#"
1802 {"time": 9223372036854775807 }
1803 {"time": -9223372036854775808 }
1804 {"time": 9e5 }
1805 "#;
1806
1807 let schema = Arc::new(Schema::new(vec![Field::new(
1808 "time",
1809 DataType::Timestamp(TimeUnit::Nanosecond, None),
1810 true,
1811 )]));
1812
1813 let batches = do_read(buf, 1024, true, false, schema);
1814 assert_eq!(batches.len(), 1);
1815
1816 let i64 = batches[0]
1817 .column(0)
1818 .as_primitive::<TimestampNanosecondType>();
1819 assert_eq!(i64.values(), &[i64::MAX, i64::MIN, 900000]);
1820 }
1821
1822 #[test]
1823 fn test_strict_mode_no_missing_columns_in_schema() {
1824 let buf = r#"
1825 {"a": 1, "b": "2", "c": true}
1826 {"a": 2E0, "b": "4", "c": false}
1827 "#;
1828
1829 let schema = Arc::new(Schema::new(vec![
1830 Field::new("a", DataType::Int16, false),
1831 Field::new("b", DataType::Utf8, false),
1832 Field::new("c", DataType::Boolean, false),
1833 ]));
1834
1835 let batches = do_read(buf, 1024, true, true, schema);
1836 assert_eq!(batches.len(), 1);
1837
1838 let buf = r#"
1839 {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1840 {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1841 "#;
1842
1843 let schema = Arc::new(Schema::new(vec![
1844 Field::new("a", DataType::Int16, false),
1845 Field::new("b", DataType::Utf8, false),
1846 Field::new_struct(
1847 "c",
1848 vec![
1849 Field::new("a", DataType::Boolean, false),
1850 Field::new("b", DataType::Int16, false),
1851 ],
1852 false,
1853 ),
1854 ]));
1855
1856 let batches = do_read(buf, 1024, true, true, schema);
1857 assert_eq!(batches.len(), 1);
1858 }
1859
1860 #[test]
1861 fn test_strict_mode_missing_columns_in_schema() {
1862 let buf = r#"
1863 {"a": 1, "b": "2", "c": true}
1864 {"a": 2E0, "b": "4", "c": false}
1865 "#;
1866
1867 let schema = Arc::new(Schema::new(vec![
1868 Field::new("a", DataType::Int16, true),
1869 Field::new("c", DataType::Boolean, true),
1870 ]));
1871
1872 let err = ReaderBuilder::new(schema)
1873 .with_batch_size(1024)
1874 .with_strict_mode(true)
1875 .build(Cursor::new(buf.as_bytes()))
1876 .unwrap()
1877 .read()
1878 .unwrap_err();
1879
1880 assert_eq!(
1881 err.to_string(),
1882 "Json error: column 'b' missing from schema"
1883 );
1884
1885 let buf = r#"
1886 {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1887 {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1888 "#;
1889
1890 let schema = Arc::new(Schema::new(vec![
1891 Field::new("a", DataType::Int16, false),
1892 Field::new("b", DataType::Utf8, false),
1893 Field::new_struct("c", vec![Field::new("a", DataType::Boolean, false)], false),
1894 ]));
1895
1896 let err = ReaderBuilder::new(schema)
1897 .with_batch_size(1024)
1898 .with_strict_mode(true)
1899 .build(Cursor::new(buf.as_bytes()))
1900 .unwrap()
1901 .read()
1902 .unwrap_err();
1903
1904 assert_eq!(
1905 err.to_string(),
1906 "Json error: whilst decoding field 'c': column 'b' missing from schema"
1907 );
1908 }
1909
1910 fn read_file(path: &str, schema: Option<Schema>) -> Reader<BufReader<File>> {
1911 let file = File::open(path).unwrap();
1912 let mut reader = BufReader::new(file);
1913 let schema = schema.unwrap_or_else(|| {
1914 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
1915 reader.rewind().unwrap();
1916 schema
1917 });
1918 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
1919 builder.build(reader).unwrap()
1920 }
1921
1922 #[test]
1923 fn test_json_basic() {
1924 let mut reader = read_file("test/data/basic.json", None);
1925 let batch = reader.next().unwrap().unwrap();
1926
1927 assert_eq!(8, batch.num_columns());
1928 assert_eq!(12, batch.num_rows());
1929
1930 let schema = reader.schema();
1931 let batch_schema = batch.schema();
1932 assert_eq!(schema, batch_schema);
1933
1934 let a = schema.column_with_name("a").unwrap();
1935 assert_eq!(0, a.0);
1936 assert_eq!(&DataType::Int64, a.1.data_type());
1937 let b = schema.column_with_name("b").unwrap();
1938 assert_eq!(1, b.0);
1939 assert_eq!(&DataType::Float64, b.1.data_type());
1940 let c = schema.column_with_name("c").unwrap();
1941 assert_eq!(2, c.0);
1942 assert_eq!(&DataType::Boolean, c.1.data_type());
1943 let d = schema.column_with_name("d").unwrap();
1944 assert_eq!(3, d.0);
1945 assert_eq!(&DataType::Utf8, d.1.data_type());
1946
1947 let aa = batch.column(a.0).as_primitive::<Int64Type>();
1948 assert_eq!(1, aa.value(0));
1949 assert_eq!(-10, aa.value(1));
1950 let bb = batch.column(b.0).as_primitive::<Float64Type>();
1951 assert_eq!(2.0, bb.value(0));
1952 assert_eq!(-3.5, bb.value(1));
1953 let cc = batch.column(c.0).as_boolean();
1954 assert!(!cc.value(0));
1955 assert!(cc.value(10));
1956 let dd = batch.column(d.0).as_string::<i32>();
1957 assert_eq!("4", dd.value(0));
1958 assert_eq!("text", dd.value(8));
1959 }
1960
1961 #[test]
1962 fn test_json_empty_projection() {
1963 let mut reader = read_file("test/data/basic.json", Some(Schema::empty()));
1964 let batch = reader.next().unwrap().unwrap();
1965
1966 assert_eq!(0, batch.num_columns());
1967 assert_eq!(12, batch.num_rows());
1968 }
1969
1970 #[test]
1971 fn test_json_basic_with_nulls() {
1972 let mut reader = read_file("test/data/basic_nulls.json", None);
1973 let batch = reader.next().unwrap().unwrap();
1974
1975 assert_eq!(4, batch.num_columns());
1976 assert_eq!(12, batch.num_rows());
1977
1978 let schema = reader.schema();
1979 let batch_schema = batch.schema();
1980 assert_eq!(schema, batch_schema);
1981
1982 let a = schema.column_with_name("a").unwrap();
1983 assert_eq!(&DataType::Int64, a.1.data_type());
1984 let b = schema.column_with_name("b").unwrap();
1985 assert_eq!(&DataType::Float64, b.1.data_type());
1986 let c = schema.column_with_name("c").unwrap();
1987 assert_eq!(&DataType::Boolean, c.1.data_type());
1988 let d = schema.column_with_name("d").unwrap();
1989 assert_eq!(&DataType::Utf8, d.1.data_type());
1990
1991 let aa = batch.column(a.0).as_primitive::<Int64Type>();
1992 assert!(aa.is_valid(0));
1993 assert!(!aa.is_valid(1));
1994 assert!(!aa.is_valid(11));
1995 let bb = batch.column(b.0).as_primitive::<Float64Type>();
1996 assert!(bb.is_valid(0));
1997 assert!(!bb.is_valid(2));
1998 assert!(!bb.is_valid(11));
1999 let cc = batch.column(c.0).as_boolean();
2000 assert!(cc.is_valid(0));
2001 assert!(!cc.is_valid(4));
2002 assert!(!cc.is_valid(11));
2003 let dd = batch.column(d.0).as_string::<i32>();
2004 assert!(!dd.is_valid(0));
2005 assert!(dd.is_valid(1));
2006 assert!(!dd.is_valid(4));
2007 assert!(!dd.is_valid(11));
2008 }
2009
2010 #[test]
2011 fn test_json_basic_schema() {
2012 let schema = Schema::new(vec![
2013 Field::new("a", DataType::Int64, true),
2014 Field::new("b", DataType::Float32, false),
2015 Field::new("c", DataType::Boolean, false),
2016 Field::new("d", DataType::Utf8, false),
2017 ]);
2018
2019 let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2020 let reader_schema = reader.schema();
2021 assert_eq!(reader_schema.as_ref(), &schema);
2022 let batch = reader.next().unwrap().unwrap();
2023
2024 assert_eq!(4, batch.num_columns());
2025 assert_eq!(12, batch.num_rows());
2026
2027 let schema = batch.schema();
2028
2029 let a = schema.column_with_name("a").unwrap();
2030 assert_eq!(&DataType::Int64, a.1.data_type());
2031 let b = schema.column_with_name("b").unwrap();
2032 assert_eq!(&DataType::Float32, b.1.data_type());
2033 let c = schema.column_with_name("c").unwrap();
2034 assert_eq!(&DataType::Boolean, c.1.data_type());
2035 let d = schema.column_with_name("d").unwrap();
2036 assert_eq!(&DataType::Utf8, d.1.data_type());
2037
2038 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2039 assert_eq!(1, aa.value(0));
2040 assert_eq!(100000000000000, aa.value(11));
2041 let bb = batch.column(b.0).as_primitive::<Float32Type>();
2042 assert_eq!(2.0, bb.value(0));
2043 assert_eq!(-3.5, bb.value(1));
2044 }
2045
2046 #[test]
2047 fn test_json_basic_schema_projection() {
2048 let schema = Schema::new(vec![
2049 Field::new("a", DataType::Int64, true),
2050 Field::new("c", DataType::Boolean, false),
2051 ]);
2052
2053 let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2054 let batch = reader.next().unwrap().unwrap();
2055
2056 assert_eq!(2, batch.num_columns());
2057 assert_eq!(2, batch.schema().fields().len());
2058 assert_eq!(12, batch.num_rows());
2059
2060 assert_eq!(batch.schema().as_ref(), &schema);
2061
2062 let a = schema.column_with_name("a").unwrap();
2063 assert_eq!(0, a.0);
2064 assert_eq!(&DataType::Int64, a.1.data_type());
2065 let c = schema.column_with_name("c").unwrap();
2066 assert_eq!(1, c.0);
2067 assert_eq!(&DataType::Boolean, c.1.data_type());
2068 }
2069
2070 #[test]
2071 fn test_json_arrays() {
2072 let mut reader = read_file("test/data/arrays.json", None);
2073 let batch = reader.next().unwrap().unwrap();
2074
2075 assert_eq!(4, batch.num_columns());
2076 assert_eq!(3, batch.num_rows());
2077
2078 let schema = batch.schema();
2079
2080 let a = schema.column_with_name("a").unwrap();
2081 assert_eq!(&DataType::Int64, a.1.data_type());
2082 let b = schema.column_with_name("b").unwrap();
2083 assert_eq!(
2084 &DataType::List(Arc::new(Field::new_list_field(DataType::Float64, true))),
2085 b.1.data_type()
2086 );
2087 let c = schema.column_with_name("c").unwrap();
2088 assert_eq!(
2089 &DataType::List(Arc::new(Field::new_list_field(DataType::Boolean, true))),
2090 c.1.data_type()
2091 );
2092 let d = schema.column_with_name("d").unwrap();
2093 assert_eq!(&DataType::Utf8, d.1.data_type());
2094
2095 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2096 assert_eq!(1, aa.value(0));
2097 assert_eq!(-10, aa.value(1));
2098 assert_eq!(1627668684594000000, aa.value(2));
2099 let bb = batch.column(b.0).as_list::<i32>();
2100 let bb = bb.values().as_primitive::<Float64Type>();
2101 assert_eq!(9, bb.len());
2102 assert_eq!(2.0, bb.value(0));
2103 assert_eq!(-6.1, bb.value(5));
2104 assert!(!bb.is_valid(7));
2105
2106 let cc = batch
2107 .column(c.0)
2108 .as_any()
2109 .downcast_ref::<ListArray>()
2110 .unwrap();
2111 let cc = cc.values().as_boolean();
2112 assert_eq!(6, cc.len());
2113 assert!(!cc.value(0));
2114 assert!(!cc.value(4));
2115 assert!(!cc.is_valid(5));
2116 }
2117
2118 #[test]
2119 fn test_empty_json_arrays() {
2120 let json_content = r#"
2121 {"items": []}
2122 {"items": null}
2123 {}
2124 "#;
2125
2126 let schema = Arc::new(Schema::new(vec![Field::new(
2127 "items",
2128 DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2129 true,
2130 )]));
2131
2132 let batches = do_read(json_content, 1024, false, false, schema);
2133 assert_eq!(batches.len(), 1);
2134
2135 let col1 = batches[0].column(0).as_list::<i32>();
2136 assert_eq!(col1.null_count(), 2);
2137 assert!(col1.value(0).is_empty());
2138 assert_eq!(col1.value(0).data_type(), &DataType::Null);
2139 assert!(col1.is_null(1));
2140 assert!(col1.is_null(2));
2141 }
2142
2143 #[test]
2144 fn test_nested_empty_json_arrays() {
2145 let json_content = r#"
2146 {"items": [[],[]]}
2147 {"items": [[null, null],[null]]}
2148 "#;
2149
2150 let schema = Arc::new(Schema::new(vec![Field::new(
2151 "items",
2152 DataType::List(FieldRef::new(Field::new_list_field(
2153 DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2154 true,
2155 ))),
2156 true,
2157 )]));
2158
2159 let batches = do_read(json_content, 1024, false, false, schema);
2160 assert_eq!(batches.len(), 1);
2161
2162 let col1 = batches[0].column(0).as_list::<i32>();
2163 assert_eq!(col1.null_count(), 0);
2164 assert_eq!(col1.value(0).len(), 2);
2165 assert!(col1.value(0).as_list::<i32>().value(0).is_empty());
2166 assert!(col1.value(0).as_list::<i32>().value(1).is_empty());
2167
2168 assert_eq!(col1.value(1).len(), 2);
2169 assert_eq!(col1.value(1).as_list::<i32>().value(0).len(), 2);
2170 assert_eq!(col1.value(1).as_list::<i32>().value(1).len(), 1);
2171 }
2172
2173 #[test]
2174 fn test_nested_list_json_arrays() {
2175 let c_field = Field::new_struct("c", vec![Field::new("d", DataType::Utf8, true)], true);
2176 let a_struct_field = Field::new_struct(
2177 "a",
2178 vec![Field::new("b", DataType::Boolean, true), c_field.clone()],
2179 true,
2180 );
2181 let a_field = Field::new("a", DataType::List(Arc::new(a_struct_field.clone())), true);
2182 let schema = Arc::new(Schema::new(vec![a_field.clone()]));
2183 let builder = ReaderBuilder::new(schema).with_batch_size(64);
2184 let json_content = r#"
2185 {"a": [{"b": true, "c": {"d": "a_text"}}, {"b": false, "c": {"d": "b_text"}}]}
2186 {"a": [{"b": false, "c": null}]}
2187 {"a": [{"b": true, "c": {"d": "c_text"}}, {"b": null, "c": {"d": "d_text"}}, {"b": true, "c": {"d": null}}]}
2188 {"a": null}
2189 {"a": []}
2190 {"a": [null]}
2191 "#;
2192 let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2193
2194 let d = StringArray::from(vec![
2196 Some("a_text"),
2197 Some("b_text"),
2198 None,
2199 Some("c_text"),
2200 Some("d_text"),
2201 None,
2202 None,
2203 ]);
2204 let c = StructArray::new(
2205 vec![Field::new("d", DataType::Utf8, true)].into(),
2206 vec![Arc::new(d.clone()) as ArrayRef],
2207 Some(NullBuffer::from(vec![
2208 true, true, false, true, true, true, false,
2209 ])),
2210 );
2211 let b = BooleanArray::from(vec![
2212 Some(true),
2213 Some(false),
2214 Some(false),
2215 Some(true),
2216 None,
2217 Some(true),
2218 None,
2219 ]);
2220 let a = StructArray::new(
2221 vec![Field::new("b", DataType::Boolean, true), c_field.clone()].into(),
2222 vec![
2223 Arc::new(b.clone()) as ArrayRef,
2224 Arc::new(c.clone()) as ArrayRef,
2225 ],
2226 Some(NullBuffer::from(vec![
2227 true, true, true, true, true, true, false,
2228 ])),
2229 );
2230 let a_list = ListArray::new(
2231 Arc::new(a_struct_field.clone()),
2232 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 6, 6, 6, 7])),
2233 Arc::new(a),
2234 Some(NullBuffer::from(vec![true, true, true, false, true, true])),
2235 );
2236
2237 let batch = reader.next().unwrap().unwrap();
2239 let read = batch.column(0);
2240 assert_eq!(read.len(), 6);
2241 let read: &ListArray = read.as_list::<i32>();
2243 let expected = &a_list;
2244 assert_eq!(read.value_offsets(), &[0, 2, 3, 6, 6, 6, 7]);
2245 assert_eq!(read.nulls(), expected.nulls());
2247 let struct_array = read.values().as_struct();
2249 let expected_struct_array = expected.values().as_struct();
2250
2251 assert_eq!(7, struct_array.len());
2252 assert_eq!(1, struct_array.null_count());
2253 assert_eq!(7, expected_struct_array.len());
2254 assert_eq!(1, expected_struct_array.null_count());
2255 assert_eq!(struct_array.nulls(), expected_struct_array.nulls());
2257 let read_b = struct_array.column(0);
2259 assert_eq!(read_b.as_ref(), &b);
2260 let read_c = struct_array.column(1);
2261 assert_eq!(read_c.as_struct(), &c);
2262 let read_c = read_c.as_struct();
2263 let read_d = read_c.column(0);
2264 assert_eq!(read_d.as_ref(), &d);
2265
2266 assert_eq!(read, expected);
2267 }
2268
2269 fn assert_read_list_view<O: OffsetSizeTrait>() {
2270 let field = Arc::new(Field::new("item", DataType::Int32, true));
2271 let data_type = GenericListViewArray::<O>::DATA_TYPE_CONSTRUCTOR(field.clone());
2272 let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2273
2274 let buf = r#"
2275 {"lv": [1, 2, 3]}
2276 {"lv": [4, null]}
2277 {"lv": null}
2278 {"lv": [6]}
2279 {"lv": []}
2280 "#;
2281
2282 let batches = do_read(buf, 1024, false, false, schema);
2283 assert_eq!(batches.len(), 1);
2284 let batch = &batches[0];
2285 let col = batch.column(0);
2286 let list_view = col
2287 .as_any()
2288 .downcast_ref::<GenericListViewArray<O>>()
2289 .unwrap();
2290
2291 assert_eq!(list_view.len(), 5);
2292
2293 let expected_offsets: Vec<O> = vec![0, 3, 5, 5, 6]
2295 .into_iter()
2296 .map(|v| O::usize_as(v))
2297 .collect();
2298 let expected_sizes: Vec<O> = vec![3, 2, 0, 1, 0]
2299 .into_iter()
2300 .map(|v| O::usize_as(v))
2301 .collect();
2302 assert_eq!(list_view.value_offsets(), &expected_offsets);
2303 assert_eq!(list_view.value_sizes(), &expected_sizes);
2304
2305 assert!(list_view.is_valid(0));
2307 let vals = list_view.value(0);
2308 let ints = vals.as_primitive::<Int32Type>();
2309 assert_eq!(ints.values(), &[1, 2, 3]);
2310
2311 assert!(list_view.is_valid(1));
2313 let vals = list_view.value(1);
2314 let ints = vals.as_primitive::<Int32Type>();
2315 assert_eq!(ints.len(), 2);
2316 assert_eq!(ints.value(0), 4);
2317 assert!(ints.is_null(1));
2318
2319 assert!(list_view.is_null(2));
2321
2322 assert!(list_view.is_valid(3));
2324 let vals = list_view.value(3);
2325 let ints = vals.as_primitive::<Int32Type>();
2326 assert_eq!(ints.values(), &[6]);
2327
2328 assert!(list_view.is_valid(4));
2330 let vals = list_view.value(4);
2331 assert_eq!(vals.len(), 0);
2332 }
2333
2334 #[test]
2335 fn test_read_list_view() {
2336 assert_read_list_view::<i32>();
2337 assert_read_list_view::<i64>();
2338 }
2339
2340 #[test]
2341 fn test_read_list_view_rejects_null_non_nullable_child() {
2342 let field = Arc::new(Field::new("item", DataType::Int32, false));
2343 for (data_type, array_type) in [
2344 (DataType::ListView(field.clone()), "ListViewArray"),
2345 (DataType::LargeListView(field.clone()), "LargeListViewArray"),
2346 ] {
2347 let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2348 let buf = r#"
2349 {"lv": [1, 2, 3]}
2350 {"lv": [4, null]}
2351 "#;
2352
2353 let error = ReaderBuilder::new(schema)
2354 .build(Cursor::new(buf.as_bytes()))
2355 .unwrap()
2356 .collect::<Result<Vec<_>, _>>()
2357 .unwrap_err();
2358
2359 assert_eq!(
2360 error.to_string(),
2361 format!(
2362 "Invalid argument error: Non-nullable field of {array_type} \"item\" cannot contain nulls"
2363 )
2364 );
2365 }
2366 }
2367
2368 #[test]
2369 fn test_fixed_size_list() {
2370 let buf = r#"
2371 {"a": [1, 2, 3]}
2372 {"a": [4, 5, 6]}
2373 {"a": [7, 8, 9]}
2374 "#;
2375
2376 let field = Field::new_list_field(DataType::Int32, true);
2377 let schema = Arc::new(Schema::new(vec![Field::new(
2378 "a",
2379 DataType::FixedSizeList(Arc::new(field), 3),
2380 false,
2381 )]));
2382
2383 let batches = do_read(buf, 1024, false, false, schema);
2384 assert_eq!(batches.len(), 1);
2385
2386 let col = batches[0].column(0).as_fixed_size_list();
2387 assert_eq!(col.len(), 3);
2388 assert_eq!(col.value_length(), 3);
2389
2390 let values = col.values().as_primitive::<Int32Type>();
2391 assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8, 9]);
2392 }
2393
2394 #[test]
2395 fn test_fixed_size_list_nullable() {
2396 let buf = r#"
2397 {"a": [1, 2]}
2398 {"a": null}
2399 {"a": [3, null]}
2400 "#;
2401
2402 let field = Field::new_list_field(DataType::Int32, true);
2403 let schema = Arc::new(Schema::new(vec![Field::new(
2404 "a",
2405 DataType::FixedSizeList(Arc::new(field), 2),
2406 true,
2407 )]));
2408
2409 let batches = do_read(buf, 1024, false, false, schema);
2410 assert_eq!(batches.len(), 1);
2411
2412 let col = batches[0].column(0).as_fixed_size_list();
2413 assert_eq!(col.len(), 3);
2414 assert!(col.is_valid(0));
2415 assert!(col.is_null(1));
2416 assert!(col.is_valid(2));
2417
2418 let values = col.values().as_primitive::<Int32Type>();
2419 assert_eq!(values.value(0), 1);
2420 assert_eq!(values.value(1), 2);
2421 assert_eq!(values.value(4), 3);
2422 assert!(values.is_null(5));
2423 }
2424
2425 #[test]
2426 fn test_fixed_size_list_zero_size_non_nullable() {
2427 let buf = r#"
2428 {"a": []}
2429 {"a": []}
2430 {"a": []}
2431 "#;
2432
2433 let field = Field::new_list_field(DataType::Int32, true);
2434 let schema = Arc::new(Schema::new(vec![Field::new(
2435 "a",
2436 DataType::FixedSizeList(Arc::new(field), 0),
2437 false,
2438 )]));
2439
2440 let batches = do_read(buf, 1024, false, false, schema);
2441 assert_eq!(batches.len(), 1);
2442
2443 let col = batches[0].column(0).as_fixed_size_list();
2444 assert_eq!(col.len(), 3);
2445 assert_eq!(col.value_length(), 0);
2446
2447 let values = col.values().as_primitive::<Int32Type>();
2448 assert!(values.values().is_empty());
2449 }
2450
2451 #[test]
2452 fn test_fixed_size_list_wrong_size() {
2453 let buf = r#"{"a": [1, 2, 3]}"#;
2454
2455 let field = Field::new_list_field(DataType::Int32, true);
2456 let schema = Arc::new(Schema::new(vec![Field::new(
2457 "a",
2458 DataType::FixedSizeList(Arc::new(field), 2),
2459 false,
2460 )]));
2461
2462 let err = ReaderBuilder::new(schema)
2463 .build(Cursor::new(buf.as_bytes()))
2464 .unwrap()
2465 .next()
2466 .unwrap()
2467 .unwrap_err();
2468
2469 assert!(err.to_string().contains("expected 2 but got 3"), "{}", err);
2470 }
2471
2472 #[test]
2473 fn test_fixed_size_list_nested() {
2474 let buf = r#"
2475 {"a": [[1, 2], [3, 4]]}
2476 {"a": [[5, 6], [7, 8]]}
2477 "#;
2478
2479 let inner_field = Field::new_list_field(DataType::Int32, true);
2480 let inner_type = DataType::FixedSizeList(Arc::new(inner_field), 2);
2481 let outer_field = Arc::new(Field::new_list_field(inner_type.clone(), true));
2482 let schema = Arc::new(Schema::new(vec![Field::new(
2483 "a",
2484 DataType::FixedSizeList(outer_field, 2),
2485 false,
2486 )]));
2487
2488 let batches = do_read(buf, 1024, false, false, schema);
2489 assert_eq!(batches.len(), 1);
2490
2491 let col = batches[0].column(0).as_fixed_size_list();
2492 assert_eq!(col.len(), 2);
2493 assert_eq!(col.value_length(), 2);
2494
2495 let inner = col.values().as_fixed_size_list();
2496 assert_eq!(inner.len(), 4);
2497 assert_eq!(inner.value_length(), 2);
2498
2499 let values = inner.values().as_primitive::<Int32Type>();
2500 assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8]);
2501 }
2502
2503 #[test]
2504 fn test_fixed_size_list_ignore_type_conflicts() {
2505 let field = Field::new("item", DataType::Int32, true);
2506 let schema = Arc::new(Schema::new(vec![Field::new(
2507 "a",
2508 DataType::FixedSizeList(Arc::new(field), 2),
2509 true,
2510 )]));
2511
2512 let json = vec![
2513 json!({"a": [1, 2]}),
2514 json!({"a": "not a list"}),
2515 json!({"a": 42}),
2516 json!({"a": [6, 7]}),
2517 ];
2518
2519 let mut decoder = ReaderBuilder::new(schema)
2520 .with_ignore_type_conflicts(true)
2521 .build_decoder()
2522 .unwrap();
2523 decoder.serialize(&json).unwrap();
2524 let batch = decoder.flush().unwrap().unwrap();
2525
2526 let col = batch.column(0).as_fixed_size_list();
2527 assert_eq!(col.len(), 4);
2528 assert!(col.is_valid(0));
2529 assert!(col.is_null(1)); assert!(col.is_null(2)); assert!(col.is_valid(3));
2532
2533 let values = col.values().as_primitive::<Int32Type>();
2534 assert_eq!(values.value(0), 1);
2535 assert_eq!(values.value(1), 2);
2536 assert_eq!(values.value(6), 6);
2537 assert_eq!(values.value(7), 7);
2538 }
2539
2540 #[test]
2541 fn test_skip_empty_lines() {
2542 let schema = Schema::new(vec![Field::new("a", DataType::Int64, true)]);
2543 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
2544 let json_content = "
2545 {\"a\": 1}
2546 {\"a\": 2}
2547 {\"a\": 3}";
2548 let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2549 let batch = reader.next().unwrap().unwrap();
2550
2551 assert_eq!(1, batch.num_columns());
2552 assert_eq!(3, batch.num_rows());
2553
2554 let schema = reader.schema();
2555 let c = schema.column_with_name("a").unwrap();
2556 assert_eq!(&DataType::Int64, c.1.data_type());
2557 }
2558
2559 #[test]
2560 fn test_with_multiple_batches() {
2561 let file = File::open("test/data/basic_nulls.json").unwrap();
2562 let mut reader = BufReader::new(file);
2563 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2564 reader.rewind().unwrap();
2565
2566 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2567 let mut reader = builder.build(reader).unwrap();
2568
2569 let mut num_records = Vec::new();
2570 while let Some(rb) = reader.next().transpose().unwrap() {
2571 num_records.push(rb.num_rows());
2572 }
2573
2574 assert_eq!(vec![5, 5, 2], num_records);
2575 }
2576
2577 #[test]
2578 fn test_timestamp_from_json_seconds() {
2579 let schema = Schema::new(vec![Field::new(
2580 "a",
2581 DataType::Timestamp(TimeUnit::Second, None),
2582 true,
2583 )]);
2584
2585 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2586 let batch = reader.next().unwrap().unwrap();
2587
2588 assert_eq!(1, batch.num_columns());
2589 assert_eq!(12, batch.num_rows());
2590
2591 let schema = reader.schema();
2592 let batch_schema = batch.schema();
2593 assert_eq!(schema, batch_schema);
2594
2595 let a = schema.column_with_name("a").unwrap();
2596 assert_eq!(
2597 &DataType::Timestamp(TimeUnit::Second, None),
2598 a.1.data_type()
2599 );
2600
2601 let aa = batch.column(a.0).as_primitive::<TimestampSecondType>();
2602 assert!(aa.is_valid(0));
2603 assert!(!aa.is_valid(1));
2604 assert!(!aa.is_valid(2));
2605 assert_eq!(1, aa.value(0));
2606 assert_eq!(1, aa.value(3));
2607 assert_eq!(5, aa.value(7));
2608 }
2609
2610 #[test]
2611 fn test_timestamp_from_json_milliseconds() {
2612 let schema = Schema::new(vec![Field::new(
2613 "a",
2614 DataType::Timestamp(TimeUnit::Millisecond, None),
2615 true,
2616 )]);
2617
2618 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2619 let batch = reader.next().unwrap().unwrap();
2620
2621 assert_eq!(1, batch.num_columns());
2622 assert_eq!(12, batch.num_rows());
2623
2624 let schema = reader.schema();
2625 let batch_schema = batch.schema();
2626 assert_eq!(schema, batch_schema);
2627
2628 let a = schema.column_with_name("a").unwrap();
2629 assert_eq!(
2630 &DataType::Timestamp(TimeUnit::Millisecond, None),
2631 a.1.data_type()
2632 );
2633
2634 let aa = batch.column(a.0).as_primitive::<TimestampMillisecondType>();
2635 assert!(aa.is_valid(0));
2636 assert!(!aa.is_valid(1));
2637 assert!(!aa.is_valid(2));
2638 assert_eq!(1, aa.value(0));
2639 assert_eq!(1, aa.value(3));
2640 assert_eq!(5, aa.value(7));
2641 }
2642
2643 #[test]
2644 fn test_date_from_json_milliseconds() {
2645 let schema = Schema::new(vec![Field::new("a", DataType::Date64, true)]);
2646
2647 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2648 let batch = reader.next().unwrap().unwrap();
2649
2650 assert_eq!(1, batch.num_columns());
2651 assert_eq!(12, batch.num_rows());
2652
2653 let schema = reader.schema();
2654 let batch_schema = batch.schema();
2655 assert_eq!(schema, batch_schema);
2656
2657 let a = schema.column_with_name("a").unwrap();
2658 assert_eq!(&DataType::Date64, a.1.data_type());
2659
2660 let aa = batch.column(a.0).as_primitive::<Date64Type>();
2661 assert!(aa.is_valid(0));
2662 assert!(!aa.is_valid(1));
2663 assert!(!aa.is_valid(2));
2664 assert_eq!(1, aa.value(0));
2665 assert_eq!(1, aa.value(3));
2666 assert_eq!(5, aa.value(7));
2667 }
2668
2669 #[test]
2670 fn test_time_from_json_nanoseconds() {
2671 let schema = Schema::new(vec![Field::new(
2672 "a",
2673 DataType::Time64(TimeUnit::Nanosecond),
2674 true,
2675 )]);
2676
2677 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2678 let batch = reader.next().unwrap().unwrap();
2679
2680 assert_eq!(1, batch.num_columns());
2681 assert_eq!(12, batch.num_rows());
2682
2683 let schema = reader.schema();
2684 let batch_schema = batch.schema();
2685 assert_eq!(schema, batch_schema);
2686
2687 let a = schema.column_with_name("a").unwrap();
2688 assert_eq!(&DataType::Time64(TimeUnit::Nanosecond), a.1.data_type());
2689
2690 let aa = batch.column(a.0).as_primitive::<Time64NanosecondType>();
2691 assert!(aa.is_valid(0));
2692 assert!(!aa.is_valid(1));
2693 assert!(!aa.is_valid(2));
2694 assert_eq!(1, aa.value(0));
2695 assert_eq!(1, aa.value(3));
2696 assert_eq!(5, aa.value(7));
2697 }
2698
2699 #[test]
2700 fn test_json_iterator() {
2701 let file = File::open("test/data/basic.json").unwrap();
2702 let mut reader = BufReader::new(file);
2703 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2704 reader.rewind().unwrap();
2705
2706 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2707 let reader = builder.build(reader).unwrap();
2708 let schema = reader.schema();
2709 let (col_a_index, _) = schema.column_with_name("a").unwrap();
2710
2711 let mut sum_num_rows = 0;
2712 let mut num_batches = 0;
2713 let mut sum_a = 0;
2714 for batch in reader {
2715 let batch = batch.unwrap();
2716 assert_eq!(8, batch.num_columns());
2717 sum_num_rows += batch.num_rows();
2718 num_batches += 1;
2719 let batch_schema = batch.schema();
2720 assert_eq!(schema, batch_schema);
2721 let a_array = batch.column(col_a_index).as_primitive::<Int64Type>();
2722 sum_a += (0..a_array.len()).map(|i| a_array.value(i)).sum::<i64>();
2723 }
2724 assert_eq!(12, sum_num_rows);
2725 assert_eq!(3, num_batches);
2726 assert_eq!(100000000000011, sum_a);
2727 }
2728
2729 #[test]
2730 fn test_decoder_error() {
2731 let schema = Arc::new(Schema::new(vec![Field::new_struct(
2732 "a",
2733 vec![Field::new("child", DataType::Int32, false)],
2734 true,
2735 )]));
2736
2737 let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
2738 let _ = decoder.decode(r#"{"a": { "child":"#.as_bytes()).unwrap();
2739 assert!(decoder.tape_decoder.has_partial_row());
2740 assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2741 let _ = decoder.flush().unwrap_err();
2742 assert!(decoder.tape_decoder.has_partial_row());
2743 assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2744
2745 let parse_err = |s: &str| {
2746 ReaderBuilder::new(schema.clone())
2747 .build(Cursor::new(s.as_bytes()))
2748 .unwrap()
2749 .next()
2750 .unwrap()
2751 .unwrap_err()
2752 .to_string()
2753 };
2754
2755 let err = parse_err(r#"{"a": 123}"#);
2756 assert_eq!(
2757 err,
2758 "Json error: whilst decoding field 'a': expected { got 123"
2759 );
2760
2761 let err = parse_err(r#"{"a": ["bar"]}"#);
2762 assert_eq!(
2763 err,
2764 r#"Json error: whilst decoding field 'a': expected { got ["bar"]"#
2765 );
2766
2767 let err = parse_err(r#"{"a": []}"#);
2768 assert_eq!(
2769 err,
2770 "Json error: whilst decoding field 'a': expected { got []"
2771 );
2772
2773 let err = parse_err(r#"{"a": [{"child": 234}]}"#);
2774 assert_eq!(
2775 err,
2776 r#"Json error: whilst decoding field 'a': expected { got [{"child": 234}]"#
2777 );
2778
2779 let err = parse_err(r#"{"a": [{"child": {"foo": [{"foo": ["bar"]}]}}]}"#);
2780 assert_eq!(
2781 err,
2782 r#"Json error: whilst decoding field 'a': expected { got [{"child": {"foo": [{"foo": ["bar"]}]}}]"#
2783 );
2784
2785 let err = parse_err(r#"{"a": true}"#);
2786 assert_eq!(
2787 err,
2788 "Json error: whilst decoding field 'a': expected { got true"
2789 );
2790
2791 let err = parse_err(r#"{"a": false}"#);
2792 assert_eq!(
2793 err,
2794 "Json error: whilst decoding field 'a': expected { got false"
2795 );
2796
2797 let err = parse_err(r#"{"a": "foo"}"#);
2798 assert_eq!(
2799 err,
2800 "Json error: whilst decoding field 'a': expected { got \"foo\""
2801 );
2802
2803 let err = parse_err(r#"{"a": {"child": false}}"#);
2804 assert_eq!(
2805 err,
2806 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got false"
2807 );
2808
2809 let err = parse_err(r#"{"a": {"child": []}}"#);
2810 assert_eq!(
2811 err,
2812 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got []"
2813 );
2814
2815 let err = parse_err(r#"{"a": {"child": [123]}}"#);
2816 assert_eq!(
2817 err,
2818 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123]"
2819 );
2820
2821 let err = parse_err(r#"{"a": {"child": [123, 3465346]}}"#);
2822 assert_eq!(
2823 err,
2824 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123, 3465346]"
2825 );
2826 }
2827
2828 #[test]
2829 fn test_serialize_timestamp() {
2830 let json = vec![
2831 json!({"timestamp": 1681319393}),
2832 json!({"timestamp": "1970-01-01T00:00:00+02:00"}),
2833 ];
2834 let schema = Schema::new(vec![Field::new(
2835 "timestamp",
2836 DataType::Timestamp(TimeUnit::Second, None),
2837 true,
2838 )]);
2839 let mut decoder = ReaderBuilder::new(Arc::new(schema))
2840 .build_decoder()
2841 .unwrap();
2842 decoder.serialize(&json).unwrap();
2843 let batch = decoder.flush().unwrap().unwrap();
2844 assert_eq!(batch.num_rows(), 2);
2845 assert_eq!(batch.num_columns(), 1);
2846 let values = batch.column(0).as_primitive::<TimestampSecondType>();
2847 assert_eq!(values.values(), &[1681319393, -7200]);
2848 }
2849
2850 #[test]
2851 fn test_serialize_decimal() {
2852 let json = vec![
2853 json!({"decimal": 1.234}),
2854 json!({"decimal": "1.234"}),
2855 json!({"decimal": 1234}),
2856 json!({"decimal": "1234"}),
2857 ];
2858 let schema = Schema::new(vec![Field::new(
2859 "decimal",
2860 DataType::Decimal128(10, 3),
2861 true,
2862 )]);
2863 let mut decoder = ReaderBuilder::new(Arc::new(schema))
2864 .build_decoder()
2865 .unwrap();
2866 decoder.serialize(&json).unwrap();
2867 let batch = decoder.flush().unwrap().unwrap();
2868 assert_eq!(batch.num_rows(), 4);
2869 assert_eq!(batch.num_columns(), 1);
2870 let values = batch.column(0).as_primitive::<Decimal128Type>();
2871 assert_eq!(values.values(), &[1234, 1234, 1234000, 1234000]);
2872 }
2873
2874 #[test]
2875 fn test_serde_field() {
2876 let field = Field::new("int", DataType::Int32, true);
2877 let mut decoder = ReaderBuilder::new_with_field(field)
2878 .build_decoder()
2879 .unwrap();
2880 decoder.serialize(&[1_i32, 2, 3, 4]).unwrap();
2881 let b = decoder.flush().unwrap().unwrap();
2882 let values = b.column(0).as_primitive::<Int32Type>().values();
2883 assert_eq!(values, &[1, 2, 3, 4]);
2884 }
2885
2886 #[test]
2887 fn test_serde_large_numbers() {
2888 let field = Field::new("int", DataType::Int64, true);
2889 let mut decoder = ReaderBuilder::new_with_field(field)
2890 .build_decoder()
2891 .unwrap();
2892
2893 decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2894 let b = decoder.flush().unwrap().unwrap();
2895 let values = b.column(0).as_primitive::<Int64Type>().values();
2896 assert_eq!(values, &[1699148028689, 2, 3, 4]);
2897
2898 let field = Field::new(
2899 "int",
2900 DataType::Timestamp(TimeUnit::Microsecond, None),
2901 true,
2902 );
2903 let mut decoder = ReaderBuilder::new_with_field(field)
2904 .build_decoder()
2905 .unwrap();
2906
2907 decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2908 let b = decoder.flush().unwrap().unwrap();
2909 let values = b
2910 .column(0)
2911 .as_primitive::<TimestampMicrosecondType>()
2912 .values();
2913 assert_eq!(values, &[1699148028689, 2, 3, 4]);
2914 }
2915
2916 #[test]
2917 fn test_coercing_primitive_into_string_decoder() {
2918 let buf = &format!(
2919 r#"[{{"a": 1, "b": "A", "c": "T"}}, {{"a": 2, "b": "BB", "c": "F"}}, {{"a": {}, "b": 123, "c": false}}, {{"a": {}, "b": 789, "c": true}}]"#,
2920 (i32::MAX as i64 + 10),
2921 i64::MAX - 10
2922 );
2923 let schema = Schema::new(vec![
2924 Field::new("a", DataType::Float64, true),
2925 Field::new("b", DataType::Utf8, true),
2926 Field::new("c", DataType::Utf8, true),
2927 ]);
2928 let json_array: Vec<serde_json::Value> = serde_json::from_str(buf).unwrap();
2929 let schema_ref = Arc::new(schema);
2930
2931 let reader = ReaderBuilder::new(schema_ref.clone()).with_coerce_primitive(true);
2933 let mut decoder = reader.build_decoder().unwrap();
2934 decoder.serialize(json_array.as_slice()).unwrap();
2935 let batch = decoder.flush().unwrap().unwrap();
2936 assert_eq!(
2937 batch,
2938 RecordBatch::try_new(
2939 schema_ref,
2940 vec![
2941 Arc::new(Float64Array::from(vec![
2942 1.0,
2943 2.0,
2944 (i32::MAX as i64 + 10) as f64,
2945 (i64::MAX - 10) as f64
2946 ])),
2947 Arc::new(StringArray::from(vec!["A", "BB", "123", "789"])),
2948 Arc::new(StringArray::from(vec!["T", "F", "false", "true"])),
2949 ]
2950 )
2951 .unwrap()
2952 );
2953 }
2954
2955 #[test]
2956 fn test_serialize_f32_into_string() {
2957 let field = Field::new("f", DataType::Utf8, true);
2959 let mut decoder = ReaderBuilder::new_with_field(field)
2960 .with_coerce_primitive(true)
2961 .build_decoder()
2962 .unwrap();
2963 decoder.serialize(&[1.5_f32, -2.25_f32]).unwrap();
2964 let batch = decoder.flush().unwrap().unwrap();
2965 let values = batch.column(0).as_string::<i32>();
2966 assert_eq!(values.value(0), "1.5");
2967 assert_eq!(values.value(1), "-2.25");
2968 }
2969
2970 fn _parse_structs(
2975 row: &str,
2976 struct_mode: StructMode,
2977 fields: Fields,
2978 as_struct: bool,
2979 ) -> Result<RecordBatch, ArrowError> {
2980 let builder = if as_struct {
2981 ReaderBuilder::new_with_field(Field::new("r", DataType::Struct(fields), true))
2982 } else {
2983 ReaderBuilder::new(Arc::new(Schema::new(fields)))
2984 };
2985 builder
2986 .with_struct_mode(struct_mode)
2987 .build(Cursor::new(row.as_bytes()))
2988 .unwrap()
2989 .next()
2990 .unwrap()
2991 }
2992
2993 #[test]
2994 fn test_struct_decoding_list_length() {
2995 use arrow_array::array;
2996
2997 let row = "[1, 2]";
2998
2999 let mut fields = vec![Field::new("a", DataType::Int32, true)];
3000 let too_few_fields = Fields::from(fields.clone());
3001 fields.push(Field::new("b", DataType::Int32, true));
3002 let correct_fields = Fields::from(fields.clone());
3003 fields.push(Field::new("c", DataType::Int32, true));
3004 let too_many_fields = Fields::from(fields.clone());
3005
3006 let parse = |fields: Fields, as_struct: bool| {
3007 _parse_structs(row, StructMode::ListOnly, fields, as_struct)
3008 };
3009
3010 let expected_row = StructArray::new(
3011 correct_fields.clone(),
3012 vec![
3013 Arc::new(array::Int32Array::from(vec![1])),
3014 Arc::new(array::Int32Array::from(vec![2])),
3015 ],
3016 None,
3017 );
3018 let row_field = Field::new("r", DataType::Struct(correct_fields.clone()), true);
3019
3020 assert_eq!(
3021 parse(too_few_fields.clone(), true).unwrap_err().to_string(),
3022 "Json error: found extra columns for 1 fields".to_string()
3023 );
3024 assert_eq!(
3025 parse(too_few_fields, false).unwrap_err().to_string(),
3026 "Json error: found extra columns for 1 fields".to_string()
3027 );
3028 assert_eq!(
3029 parse(correct_fields.clone(), true).unwrap(),
3030 RecordBatch::try_new(
3031 Arc::new(Schema::new(vec![row_field])),
3032 vec![Arc::new(expected_row.clone())]
3033 )
3034 .unwrap()
3035 );
3036 assert_eq!(
3037 parse(correct_fields, false).unwrap(),
3038 RecordBatch::from(expected_row)
3039 );
3040 assert_eq!(
3041 parse(too_many_fields.clone(), true)
3042 .unwrap_err()
3043 .to_string(),
3044 "Json error: found 2 columns for 3 fields".to_string()
3045 );
3046 assert_eq!(
3047 parse(too_many_fields, false).unwrap_err().to_string(),
3048 "Json error: found 2 columns for 3 fields".to_string()
3049 );
3050 }
3051
3052 #[test]
3053 fn test_struct_decoding() {
3054 use arrow_array::builder;
3055
3056 let nested_object_json = r#"{"a": {"b": [1, 2], "c": {"d": 3}}}"#;
3057 let nested_list_json = r#"[[[1, 2], {"d": 3}]]"#;
3058 let nested_mixed_json = r#"{"a": [[1, 2], {"d": 3}]}"#;
3059
3060 let struct_fields = Fields::from(vec![
3061 Field::new("b", DataType::new_list(DataType::Int32, true), true),
3062 Field::new_map(
3063 "c",
3064 "entries",
3065 Field::new("keys", DataType::Utf8, false),
3066 Field::new("values", DataType::Int32, true),
3067 false,
3068 false,
3069 ),
3070 ]);
3071
3072 let list_array =
3073 ListArray::from_iter_primitive::<Int32Type, _, _>(vec![Some(vec![Some(1), Some(2)])]);
3074
3075 let map_array = {
3076 let mut map_builder = builder::MapBuilder::new(
3077 None,
3078 builder::StringBuilder::new(),
3079 builder::Int32Builder::new(),
3080 );
3081 map_builder.keys().append_value("d");
3082 map_builder.values().append_value(3);
3083 map_builder.append(true).unwrap();
3084 map_builder.finish()
3085 };
3086
3087 let struct_array = StructArray::new(
3088 struct_fields.clone(),
3089 vec![Arc::new(list_array), Arc::new(map_array)],
3090 None,
3091 );
3092
3093 let fields = Fields::from(vec![Field::new("a", DataType::Struct(struct_fields), true)]);
3094 let schema = Arc::new(Schema::new(fields.clone()));
3095 let expected = RecordBatch::try_new(schema.clone(), vec![Arc::new(struct_array)]).unwrap();
3096
3097 let parse = |row: &str, struct_mode: StructMode| {
3098 _parse_structs(row, struct_mode, fields.clone(), false)
3099 };
3100
3101 assert_eq!(
3102 parse(nested_object_json, StructMode::ObjectOnly).unwrap(),
3103 expected
3104 );
3105 assert_eq!(
3106 parse(nested_list_json, StructMode::ObjectOnly)
3107 .unwrap_err()
3108 .to_string(),
3109 "Json error: expected { got [[[1, 2], {\"d\": 3}]]".to_owned()
3110 );
3111 assert_eq!(
3112 parse(nested_mixed_json, StructMode::ObjectOnly)
3113 .unwrap_err()
3114 .to_string(),
3115 "Json error: whilst decoding field 'a': expected { got [[1, 2], {\"d\": 3}]".to_owned()
3116 );
3117
3118 assert_eq!(
3119 parse(nested_list_json, StructMode::ListOnly).unwrap(),
3120 expected
3121 );
3122 assert_eq!(
3123 parse(nested_object_json, StructMode::ListOnly)
3124 .unwrap_err()
3125 .to_string(),
3126 "Json error: expected [ got {\"a\": {\"b\": [1, 2]\"c\": {\"d\": 3}}}".to_owned()
3127 );
3128 assert_eq!(
3129 parse(nested_mixed_json, StructMode::ListOnly)
3130 .unwrap_err()
3131 .to_string(),
3132 "Json error: expected [ got {\"a\": [[1, 2], {\"d\": 3}]}".to_owned()
3133 );
3134 }
3135
3136 #[test]
3142 fn test_struct_decoding_empty_list() {
3143 let int_field = Field::new("a", DataType::Int32, true);
3144 let struct_field = Field::new(
3145 "r",
3146 DataType::Struct(Fields::from(vec![int_field.clone()])),
3147 true,
3148 );
3149
3150 let parse = |row: &str, as_struct: bool, field: Field| {
3151 _parse_structs(
3152 row,
3153 StructMode::ListOnly,
3154 Fields::from(vec![field]),
3155 as_struct,
3156 )
3157 };
3158
3159 assert_eq!(
3161 parse("[]", true, struct_field.clone())
3162 .unwrap_err()
3163 .to_string(),
3164 "Json error: found 0 columns for 1 fields".to_owned()
3165 );
3166 assert_eq!(
3167 parse("[]", false, int_field.clone())
3168 .unwrap_err()
3169 .to_string(),
3170 "Json error: found 0 columns for 1 fields".to_owned()
3171 );
3172 assert_eq!(
3173 parse("[]", false, struct_field.clone())
3174 .unwrap_err()
3175 .to_string(),
3176 "Json error: found 0 columns for 1 fields".to_owned()
3177 );
3178 assert_eq!(
3179 parse("[[]]", false, struct_field.clone())
3180 .unwrap_err()
3181 .to_string(),
3182 "Json error: whilst decoding field 'r': found 0 columns for 1 fields".to_owned()
3183 );
3184 }
3185
3186 #[test]
3187 fn test_decode_list_struct_with_wrong_types() {
3188 let int_field = Field::new("a", DataType::Int32, true);
3189 let struct_field = Field::new(
3190 "r",
3191 DataType::Struct(Fields::from(vec![int_field.clone()])),
3192 true,
3193 );
3194
3195 let parse = |row: &str, as_struct: bool, field: Field| {
3196 _parse_structs(
3197 row,
3198 StructMode::ListOnly,
3199 Fields::from(vec![field]),
3200 as_struct,
3201 )
3202 };
3203
3204 assert_eq!(
3206 parse(r#"[["a"]]"#, false, struct_field.clone())
3207 .unwrap_err()
3208 .to_string(),
3209 "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3210 );
3211 assert_eq!(
3212 parse(r#"[["a"]]"#, true, struct_field.clone())
3213 .unwrap_err()
3214 .to_string(),
3215 "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3216 );
3217 assert_eq!(
3218 parse(r#"["a"]"#, true, int_field.clone())
3219 .unwrap_err()
3220 .to_string(),
3221 "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3222 );
3223 assert_eq!(
3224 parse(r#"["a"]"#, false, int_field.clone())
3225 .unwrap_err()
3226 .to_string(),
3227 "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3228 );
3229 }
3230
3231 #[test]
3232 fn test_type_conflict_nulls() {
3233 let schema = Schema::new(vec![
3234 Field::new("null", DataType::Null, true),
3235 Field::new("bool", DataType::Boolean, true),
3236 Field::new("primitive", DataType::Int32, true),
3237 Field::new("numeric", DataType::Decimal128(10, 3), true),
3238 Field::new("string", DataType::Utf8, true),
3239 Field::new("string_view", DataType::Utf8View, true),
3240 Field::new(
3241 "timestamp",
3242 DataType::Timestamp(TimeUnit::Second, None),
3243 true,
3244 ),
3245 Field::new(
3246 "array",
3247 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3248 true,
3249 ),
3250 Field::new(
3251 "map",
3252 DataType::Map(
3253 Arc::new(Field::new(
3254 "entries",
3255 DataType::Struct(Fields::from(vec![
3256 Field::new("keys", DataType::Utf8, false),
3257 Field::new("values", DataType::Utf8, true),
3258 ])),
3259 false, )),
3261 false, ),
3263 true, ),
3265 Field::new(
3266 "struct",
3267 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3268 true,
3269 ),
3270 ]);
3271
3272 let json_values = vec![
3274 json!(null),
3275 json!(true),
3276 json!(42),
3277 json!(1.234),
3278 json!("hi"),
3279 json!("ho"),
3280 json!("1970-01-01T00:00:00+02:00"),
3281 json!([1, "ho", 3]),
3282 json!({"k": "value"}),
3283 json!({"a": 1}),
3284 ];
3285
3286 let json: Vec<_> = (0..json_values.len())
3288 .map(|i| {
3289 let pairs = json_values[i..]
3290 .iter()
3291 .chain(json_values[..i].iter())
3292 .zip(&schema.fields)
3293 .map(|(v, f)| (f.name().to_string(), v.clone()))
3294 .collect();
3295 serde_json::Value::Object(pairs)
3296 })
3297 .collect();
3298 let mut decoder = ReaderBuilder::new(Arc::new(schema))
3299 .with_ignore_type_conflicts(true)
3300 .with_coerce_primitive(true)
3301 .build_decoder()
3302 .unwrap();
3303 decoder.serialize(&json).unwrap();
3304 let batch = decoder.flush().unwrap().unwrap();
3305 assert_eq!(batch.num_rows(), 10);
3306 assert_eq!(batch.num_columns(), 10);
3307
3308 let _ = batch
3310 .column(0)
3311 .as_any()
3312 .downcast_ref::<NullArray>()
3313 .unwrap();
3314
3315 assert!(
3316 batch
3317 .column(1)
3318 .as_any()
3319 .downcast_ref::<BooleanArray>()
3320 .unwrap()
3321 .iter()
3322 .eq([
3323 Some(true),
3324 None,
3325 None,
3326 None,
3327 None,
3328 None,
3329 None,
3330 None,
3331 None,
3332 None
3333 ])
3334 );
3335
3336 assert!(batch.column(2).as_primitive::<Int32Type>().iter().eq([
3337 Some(42),
3338 Some(1),
3339 None,
3340 None,
3341 None,
3342 None,
3343 None,
3344 None,
3345 None,
3346 None
3347 ]));
3348
3349 assert!(batch.column(3).as_primitive::<Decimal128Type>().iter().eq([
3350 Some(1234),
3351 None,
3352 None,
3353 None,
3354 None,
3355 None,
3356 None,
3357 None,
3358 None,
3359 Some(42000)
3360 ]));
3361
3362 assert!(
3363 batch
3364 .column(4)
3365 .as_any()
3366 .downcast_ref::<StringArray>()
3367 .unwrap()
3368 .iter()
3369 .eq([
3370 Some("hi"),
3371 Some("ho"),
3372 Some("1970-01-01T00:00:00+02:00"),
3373 None,
3374 None,
3375 None,
3376 None,
3377 Some("true"),
3378 Some("42"),
3379 Some("1.234"),
3380 ])
3381 );
3382
3383 assert!(
3384 batch
3385 .column(5)
3386 .as_any()
3387 .downcast_ref::<StringViewArray>()
3388 .unwrap()
3389 .iter()
3390 .eq([
3391 Some("ho"),
3392 Some("1970-01-01T00:00:00+02:00"),
3393 None,
3394 None,
3395 None,
3396 None,
3397 Some("true"),
3398 Some("42"),
3399 Some("1.234"),
3400 Some("hi"),
3401 ])
3402 );
3403
3404 assert!(
3405 batch
3406 .column(6)
3407 .as_primitive::<TimestampSecondType>()
3408 .iter()
3409 .eq([
3410 Some(-7200),
3411 None,
3412 None,
3413 None,
3414 None,
3415 None,
3416 Some(42),
3417 None,
3418 None,
3419 None,
3420 ])
3421 );
3422
3423 let arrays = batch
3424 .column(7)
3425 .as_any()
3426 .downcast_ref::<ListArray>()
3427 .unwrap();
3428 assert_eq!(
3429 arrays.nulls(),
3430 Some(&NullBuffer::from(
3431 &[
3432 true, false, false, false, false, false, false, false, false, false
3433 ][..]
3434 ))
3435 );
3436 assert_eq!(arrays.offsets()[1], 3);
3437 let array_values = arrays
3438 .values()
3439 .as_any()
3440 .downcast_ref::<Int32Array>()
3441 .unwrap();
3442 assert!(array_values.iter().eq([Some(1), None, Some(3)]));
3443
3444 let maps = batch.column(8).as_any().downcast_ref::<MapArray>().unwrap();
3445 assert_eq!(
3446 maps.nulls(),
3447 Some(&NullBuffer::from(
3448 &[
3450 true, true, false, false, false, false, false, false, false, false
3451 ][..]
3452 ))
3453 );
3454 let map_keys = maps.keys().as_any().downcast_ref::<StringArray>().unwrap();
3455 assert!(map_keys.iter().eq([Some("k"), Some("a")]));
3456 let map_values = maps
3457 .values()
3458 .as_any()
3459 .downcast_ref::<StringArray>()
3460 .unwrap();
3461 assert!(map_values.iter().eq([Some("value"), Some("1")]));
3462
3463 let structs = batch
3464 .column(9)
3465 .as_any()
3466 .downcast_ref::<StructArray>()
3467 .unwrap();
3468 assert_eq!(
3469 structs.nulls(),
3470 Some(&NullBuffer::from(
3471 &[
3473 true, false, false, false, false, false, false, false, false, true
3474 ][..]
3475 ))
3476 );
3477 let struct_fields = structs
3478 .column(0)
3479 .as_any()
3480 .downcast_ref::<Int32Array>()
3481 .unwrap();
3482 assert!(struct_fields.slice(0, 2).iter().eq([Some(1), None]));
3483 }
3484
3485 #[test]
3486 fn test_type_conflict_non_nullable() {
3487 let fields = [
3488 Field::new("bool", DataType::Boolean, false),
3489 Field::new("primitive", DataType::Int32, false),
3490 Field::new("numeric", DataType::Decimal128(10, 3), false),
3491 Field::new("string", DataType::Utf8, false),
3492 Field::new("string_view", DataType::Utf8View, false),
3493 Field::new(
3494 "timestamp",
3495 DataType::Timestamp(TimeUnit::Second, None),
3496 false,
3497 ),
3498 Field::new(
3499 "array",
3500 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3501 false,
3502 ),
3503 Field::new(
3504 "fixed_size_list",
3505 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3506 false,
3507 ),
3508 Field::new(
3509 "map",
3510 DataType::Map(
3511 Arc::new(Field::new(
3512 "entries",
3513 DataType::Struct(Fields::from(vec![
3514 Field::new("keys", DataType::Utf8, false),
3515 Field::new("values", DataType::Utf8, true),
3516 ])),
3517 false, )),
3519 false, ),
3521 false, ),
3523 Field::new(
3524 "struct",
3525 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3526 false,
3527 ),
3528 ];
3529
3530 let json_values = vec![json!(true), json!({"a": 1})];
3532
3533 for field in fields {
3534 let mut decoder = ReaderBuilder::new_with_field(field)
3535 .with_ignore_type_conflicts(true)
3536 .build_decoder()
3537 .unwrap();
3538 decoder.serialize(&json_values).unwrap();
3539 decoder
3540 .flush()
3541 .expect_err("type conflict on non-nullable type");
3542 }
3543 }
3544
3545 #[test]
3546 fn test_ignore_type_conflicts_disabled() {
3547 let fields = [
3548 Field::new("null", DataType::Null, true),
3549 Field::new("bool", DataType::Boolean, true),
3550 Field::new("primitive", DataType::Int32, true),
3551 Field::new("numeric", DataType::Decimal128(10, 3), true),
3552 Field::new("string", DataType::Utf8, true),
3553 Field::new("string_view", DataType::Utf8View, true),
3554 Field::new(
3555 "timestamp",
3556 DataType::Timestamp(TimeUnit::Second, None),
3557 true,
3558 ),
3559 Field::new(
3560 "array",
3561 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3562 true,
3563 ),
3564 Field::new(
3565 "fixed_size_list",
3566 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3567 true,
3568 ),
3569 Field::new(
3570 "map",
3571 DataType::Map(
3572 Arc::new(Field::new(
3573 "entries",
3574 DataType::Struct(Fields::from(vec![
3575 Field::new("keys", DataType::Utf8, false),
3576 Field::new("values", DataType::Utf8, true),
3577 ])),
3578 false, )),
3580 false, ),
3582 true, ),
3584 Field::new(
3585 "struct",
3586 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3587 true,
3588 ),
3589 ];
3590
3591 let json_values = vec![json!(true), json!({"a": 1})];
3593
3594 for field in fields {
3595 let mut decoder = ReaderBuilder::new_with_field(field)
3596 .build_decoder()
3597 .unwrap();
3598 decoder.serialize(&json_values).unwrap();
3599 decoder
3600 .flush()
3601 .expect_err("type conflict on non-nullable type");
3602 }
3603 }
3604
3605 #[test]
3606 fn test_read_run_end_encoded() {
3607 let buf = r#"
3608 {"a": "x"}
3609 {"a": "x"}
3610 {"a": "y"}
3611 {"a": "y"}
3612 {"a": "y"}
3613 "#;
3614
3615 let ree_type = DataType::RunEndEncoded(
3616 Arc::new(Field::new("run_ends", DataType::Int32, false)),
3617 Arc::new(Field::new("values", DataType::Utf8, true)),
3618 );
3619 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3620 let batches = do_read(buf, 1024, false, false, schema);
3621 assert_eq!(batches.len(), 1);
3622
3623 let col = batches[0].column(0);
3624 let run_array = col.as_run::<arrow_array::types::Int32Type>();
3625
3626 assert_eq!(run_array.len(), 5);
3628 assert_eq!(run_array.run_ends().values(), &[2, 5]);
3629
3630 let values = run_array.values().as_string::<i32>();
3631 assert_eq!(values.len(), 2);
3632 assert_eq!(values.value(0), "x");
3633 assert_eq!(values.value(1), "y");
3634 }
3635
3636 #[test]
3637 fn test_read_run_end_encoded_consecutive_nulls() {
3638 let buf = r#"
3639 {"a": "x"}
3640 {}
3641 {}
3642 {}
3643 {"a": "y"}
3644 "#;
3645
3646 let ree_type = DataType::RunEndEncoded(
3647 Arc::new(Field::new("run_ends", DataType::Int32, false)),
3648 Arc::new(Field::new("values", DataType::Utf8, true)),
3649 );
3650 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3651 let batches = do_read(buf, 1024, false, false, schema);
3652 assert_eq!(batches.len(), 1);
3653
3654 let col = batches[0].column(0);
3655 let run_array = col.as_run::<arrow_array::types::Int32Type>();
3656
3657 assert_eq!(run_array.len(), 5);
3659 assert_eq!(run_array.run_ends().values(), &[1, 4, 5]);
3660
3661 let values = run_array.values().as_string::<i32>();
3662 assert_eq!(values.len(), 3);
3663 assert_eq!(values.value(0), "x");
3664 assert!(values.is_null(1));
3665 assert_eq!(values.value(2), "y");
3666 }
3667
3668 #[test]
3669 fn test_read_run_end_encoded_all_unique() {
3670 let buf = r#"
3671 {"a": 1}
3672 {"a": 2}
3673 {"a": 3}
3674 "#;
3675
3676 let ree_type = DataType::RunEndEncoded(
3677 Arc::new(Field::new("run_ends", DataType::Int32, false)),
3678 Arc::new(Field::new("values", DataType::Int32, true)),
3679 );
3680 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3681 let batches = do_read(buf, 1024, false, false, schema);
3682 assert_eq!(batches.len(), 1);
3683
3684 let col = batches[0].column(0);
3685 let run_array = col.as_run::<arrow_array::types::Int32Type>();
3686
3687 assert_eq!(run_array.len(), 3);
3689 assert_eq!(run_array.run_ends().values(), &[1, 2, 3]);
3690 }
3691
3692 #[test]
3693 fn test_read_run_end_encoded_int16_run_ends() {
3694 let buf = r#"
3695 {"a": "x"}
3696 {"a": "x"}
3697 {"a": "y"}
3698 "#;
3699
3700 let ree_type = DataType::RunEndEncoded(
3701 Arc::new(Field::new("run_ends", DataType::Int16, false)),
3702 Arc::new(Field::new("values", DataType::Utf8, true)),
3703 );
3704 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3705 let batches = do_read(buf, 1024, false, false, schema);
3706 assert_eq!(batches.len(), 1);
3707
3708 let col = batches[0].column(0);
3709 let run_array = col.as_run::<arrow_array::types::Int16Type>();
3710
3711 assert_eq!(run_array.len(), 3);
3712 assert_eq!(run_array.run_ends().values(), &[2i16, 3]);
3713 }
3714
3715 #[test]
3716 fn test_read_nested_run_end_encoded() {
3717 let buf = r#"
3718 {"a": "x"}
3719 {"a": "x"}
3720 {"a": "y"}
3721 "#;
3722
3723 let inner_type = DataType::RunEndEncoded(
3726 Arc::new(Field::new("run_ends", DataType::Int64, false)),
3727 Arc::new(Field::new("values", DataType::Utf8, true)),
3728 );
3729 let outer_type = DataType::RunEndEncoded(
3730 Arc::new(Field::new("run_ends", DataType::Int64, false)),
3731 Arc::new(Field::new("values", inner_type, true)),
3732 );
3733 let schema = Arc::new(Schema::new(vec![Field::new("a", outer_type, true)]));
3734 let batches = do_read(buf, 1024, false, false, schema);
3735 assert_eq!(batches.len(), 1);
3736
3737 let col = batches[0].column(0);
3738 let outer = col.as_run::<arrow_array::types::Int64Type>();
3739 assert_eq!(outer.len(), 3);
3741 assert_eq!(outer.run_ends().values(), &[2, 3]);
3742
3743 let nested = outer.values().as_run::<arrow_array::types::Int64Type>();
3744 assert_eq!(nested.len(), 2);
3746 assert_eq!(nested.run_ends().values(), &[1, 2]);
3747
3748 let nested_values = nested.values().as_string::<i32>();
3749 assert_eq!(nested_values.len(), 2);
3750 assert_eq!(nested_values.value(0), "x");
3751 assert_eq!(nested_values.value(1), "y");
3752 }
3753}