1use std::io::BufRead;
137use std::sync::Arc;
138
139use arrow_array::cast::AsArray;
140use arrow_array::timezone::Tz;
141use arrow_array::types::*;
142use arrow_array::{ArrayRef, RecordBatch, RecordBatchReader, downcast_integer};
143use arrow_schema::{ArrowError, DataType, Field, FieldRef, Schema, SchemaRef, TimeUnit};
144use chrono::Utc;
145use serde_core::Serialize;
146
147use crate::StructMode;
148use crate::reader::binary_array::{
149 BinaryArrayDecoder, BinaryViewDecoder, FixedSizeBinaryArrayDecoder,
150};
151use crate::reader::boolean_array::BooleanArrayDecoder;
152use crate::reader::decimal_array::DecimalArrayDecoder;
153use crate::reader::list_array::{
154 FixedSizeListArrayDecoder, ListArrayDecoder, ListViewArrayDecoder,
155};
156use crate::reader::map_array::MapArrayDecoder;
157use crate::reader::null_array::NullArrayDecoder;
158use crate::reader::primitive_array::PrimitiveArrayDecoder;
159use crate::reader::run_end_array::RunEndEncodedArrayDecoder;
160use crate::reader::string_array::StringArrayDecoder;
161use crate::reader::string_view_array::StringViewArrayDecoder;
162use crate::reader::struct_array::StructArrayDecoder;
163use crate::reader::tape::{TapeDecoder, TapeDecoderOptions};
164use crate::reader::timestamp_array::TimestampArrayDecoder;
165
166pub use schema::*;
167pub use tape::{Tape, TapeElement};
169pub use value_iter::ValueIter;
170
171mod binary_array;
172mod boolean_array;
173mod decimal_array;
174mod list_array;
175mod map_array;
176mod null_array;
177mod primitive_array;
178mod run_end_array;
179mod schema;
180mod serializer;
181mod string_array;
182mod string_view_array;
183mod struct_array;
184mod tape;
185mod timestamp_array;
186mod value_iter;
187
188#[derive(Clone)]
190pub struct ReaderBuilder {
191 batch_size: usize,
192 coerce_primitive: bool,
193 strict_mode: bool,
194 ignore_type_conflicts: bool,
195 is_field: bool,
196 struct_mode: StructMode,
197 flatten_top_level_arrays: bool,
198 decoder_factory: Option<Arc<dyn DecoderFactory>>,
199 schema: SchemaRef,
200}
201
202impl ReaderBuilder {
203 pub fn new(schema: SchemaRef) -> Self {
212 Self {
213 batch_size: 1024,
214 coerce_primitive: false,
215 strict_mode: false,
216 ignore_type_conflicts: false,
217 is_field: false,
218 struct_mode: Default::default(),
219 flatten_top_level_arrays: false,
220 decoder_factory: None,
221 schema,
222 }
223 }
224
225 pub fn new_with_field(field: impl Into<FieldRef>) -> Self {
256 Self {
257 batch_size: 1024,
258 coerce_primitive: false,
259 strict_mode: false,
260 ignore_type_conflicts: false,
261 is_field: true,
262 struct_mode: Default::default(),
263 flatten_top_level_arrays: false,
264 decoder_factory: None,
265 schema: Arc::new(Schema::new([field.into()])),
266 }
267 }
268
269 pub fn with_batch_size(self, batch_size: usize) -> Self {
271 Self { batch_size, ..self }
272 }
273
274 pub fn with_coerce_primitive(self, coerce_primitive: bool) -> Self {
277 Self {
278 coerce_primitive,
279 ..self
280 }
281 }
282
283 pub fn with_strict_mode(self, strict_mode: bool) -> Self {
289 Self {
290 strict_mode,
291 ..self
292 }
293 }
294
295 pub fn with_struct_mode(self, struct_mode: StructMode) -> Self {
299 Self {
300 struct_mode,
301 ..self
302 }
303 }
304
305 pub fn with_ignore_type_conflicts(self, ignore_type_conflicts: bool) -> Self {
318 Self {
319 ignore_type_conflicts,
320 ..self
321 }
322 }
323 pub fn with_flatten(self, flatten_top_level_arrays: bool) -> Self {
341 Self {
342 flatten_top_level_arrays,
343 ..self
344 }
345 }
346
347 pub fn with_decoder_factory(self, decoder_factory: Arc<dyn DecoderFactory>) -> Self {
349 Self {
350 decoder_factory: Some(decoder_factory),
351 ..self
352 }
353 }
354
355 pub fn build<R: BufRead>(self, reader: R) -> Result<Reader<R>, ArrowError> {
357 Ok(Reader {
358 reader,
359 decoder: self.build_decoder()?,
360 })
361 }
362
363 pub fn build_decoder(self) -> Result<Decoder, ArrowError> {
365 let (field, nullable) = if self.is_field {
366 let field = self.schema.fields[0].clone();
367 let nullable = field.is_nullable();
368 (field, nullable)
369 } else {
370 let data_type = DataType::Struct(self.schema.fields.clone());
372 (Arc::new(Field::new("", data_type, false)), false)
373 };
374
375 let ctx = DecoderContext {
376 coerce_primitive: self.coerce_primitive,
377 strict_mode: self.strict_mode,
378 struct_mode: self.struct_mode,
379 ignore_type_conflicts: self.ignore_type_conflicts,
380 decoder_factory: self.decoder_factory,
381 };
382 let decoder = ctx.make_decoder(&field, nullable)?;
383
384 let num_fields = self.schema.flattened_fields().len();
385
386 Ok(Decoder {
387 decoder,
388 is_field: self.is_field,
389 tape_decoder: TapeDecoder::new(TapeDecoderOptions {
390 batch_size: self.batch_size,
391 num_fields,
392 flatten_top_level_arrays: self.flatten_top_level_arrays,
393 }),
394 batch_size: self.batch_size,
395 schema: self.schema,
396 })
397 }
398}
399
400pub struct Reader<R> {
404 reader: R,
405 decoder: Decoder,
406}
407
408impl<R> std::fmt::Debug for Reader<R> {
409 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
410 f.debug_struct("Reader")
411 .field("decoder", &self.decoder)
412 .finish()
413 }
414}
415
416impl<R: BufRead> Reader<R> {
417 fn read(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
419 loop {
420 let buf = self.reader.fill_buf()?;
421 if buf.is_empty() {
422 break;
423 }
424 let read = buf.len();
425
426 let decoded = self.decoder.decode(buf)?;
427 self.reader.consume(decoded);
428 if decoded != read {
429 break;
430 }
431 }
432 self.decoder.flush()
433 }
434}
435
436impl<R: BufRead> Iterator for Reader<R> {
437 type Item = Result<RecordBatch, ArrowError>;
438
439 fn next(&mut self) -> Option<Self::Item> {
440 self.read().transpose()
441 }
442}
443
444impl<R: BufRead> RecordBatchReader for Reader<R> {
445 fn schema(&self) -> SchemaRef {
446 self.decoder.schema.clone()
447 }
448}
449
450pub struct Decoder {
491 tape_decoder: TapeDecoder,
492 decoder: Box<dyn ArrayDecoder>,
493 batch_size: usize,
494 is_field: bool,
495 schema: SchemaRef,
496}
497
498impl std::fmt::Debug for Decoder {
499 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
500 f.debug_struct("Decoder")
501 .field("schema", &self.schema)
502 .field("batch_size", &self.batch_size)
503 .finish()
504 }
505}
506
507impl Decoder {
508 pub fn decode(&mut self, buf: &[u8]) -> Result<usize, ArrowError> {
517 self.tape_decoder.decode(buf)
518 }
519
520 pub fn serialize<S: Serialize>(&mut self, rows: &[S]) -> Result<(), ArrowError> {
697 self.tape_decoder.serialize(rows)
698 }
699
700 pub fn has_partial_record(&self) -> bool {
702 self.tape_decoder.has_partial_row()
703 }
704
705 pub fn len(&self) -> usize {
707 self.tape_decoder.num_buffered_rows()
708 }
709
710 pub fn is_empty(&self) -> bool {
712 self.len() == 0
713 }
714
715 pub fn flush(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
722 let tape = self.tape_decoder.finish()?;
723
724 if tape.num_rows() == 0 {
725 return Ok(None);
726 }
727
728 let mut next_object = 1;
730 let pos = (0..tape.num_rows())
731 .map(|_| {
732 let next = tape.next(next_object, "row")?;
733 Ok(std::mem::replace(&mut next_object, next))
734 })
735 .collect::<Result<Vec<_>, ArrowError>>()?;
736
737 let decoded = self.decoder.decode(&tape, &pos)?;
738 self.tape_decoder.clear();
739
740 let batch = match self.is_field {
741 true => RecordBatch::try_new(self.schema.clone(), vec![decoded])?,
742 false => {
743 RecordBatch::from(decoded.as_struct().clone()).with_schema(self.schema.clone())?
744 }
745 };
746
747 Ok(Some(batch))
748 }
749}
750
751pub trait ArrayDecoder: Send {
755 fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError>;
761}
762
763pub trait DecoderFactory: std::fmt::Debug + Send + Sync {
898 fn make_default_decoder(
916 &self,
917 _ctx: &DecoderContext,
918 _field: &FieldRef,
919 _is_nullable: bool,
920 ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> {
921 Ok(None)
922 }
923}
924
925struct CheckedDecoder {
928 inner: Box<dyn ArrayDecoder>,
929 data_type: DataType,
930}
931
932impl ArrayDecoder for CheckedDecoder {
933 fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
934 let array = self.inner.decode(tape, pos)?;
935 if array.data_type() != &self.data_type {
936 return Err(ArrowError::JsonError(format!(
937 "custom decoder returned {} for a field of type {}",
938 array.data_type(),
939 self.data_type
940 )));
941 }
942 if array.len() != pos.len() {
943 return Err(ArrowError::JsonError(format!(
944 "custom decoder returned {} values for {} rows",
945 array.len(),
946 pos.len()
947 )));
948 }
949 Ok(array)
950 }
951}
952
953pub struct DecoderContext {
958 coerce_primitive: bool,
960 strict_mode: bool,
962 struct_mode: StructMode,
964 ignore_type_conflicts: bool,
966 decoder_factory: Option<Arc<dyn DecoderFactory>>,
968}
969
970impl DecoderContext {
971 pub fn coerce_primitive(&self) -> bool {
973 self.coerce_primitive
974 }
975
976 pub fn strict_mode(&self) -> bool {
978 self.strict_mode
979 }
980
981 pub fn struct_mode(&self) -> StructMode {
983 self.struct_mode
984 }
985
986 pub fn ignore_type_conflicts(&self) -> bool {
988 self.ignore_type_conflicts
989 }
990
991 pub fn decoder_factory(&self) -> Option<&Arc<dyn DecoderFactory>> {
993 self.decoder_factory.as_ref()
994 }
995
996 pub fn make_decoder(
1001 &self,
1002 field: &FieldRef,
1003 is_nullable: bool,
1004 ) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
1005 make_decoder(self, field, is_nullable)
1006 }
1007
1008 pub fn make_builtin_decoder(
1014 &self,
1015 field: &FieldRef,
1016 is_nullable: bool,
1017 ) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
1018 make_builtin_decoder(self, field, is_nullable)
1019 }
1020}
1021
1022fn make_decoder(
1023 ctx: &DecoderContext,
1024 field: &FieldRef,
1025 is_nullable: bool,
1026) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
1027 if let Some(factory) = ctx.decoder_factory()
1028 && let Some(decoder) = factory.make_default_decoder(ctx, field, is_nullable)?
1029 {
1030 return Ok(Box::new(CheckedDecoder {
1031 inner: decoder,
1032 data_type: field.data_type().clone(),
1033 }));
1034 }
1035
1036 make_builtin_decoder(ctx, field, is_nullable)
1037}
1038
1039fn make_builtin_decoder(
1040 ctx: &DecoderContext,
1041 field: &FieldRef,
1042 is_nullable: bool,
1043) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
1044 let data_type = field.data_type();
1045
1046 macro_rules! primitive_decoder {
1047 ($t:ty, $data_type:expr) => {
1048 Ok(Box::new(PrimitiveArrayDecoder::<$t>::new(ctx, $data_type)))
1049 };
1050 }
1051 macro_rules! timestamp_decoder {
1052 ($t:ty, $data_type:expr, $tz:expr) => {{
1053 Ok(Box::new(TimestampArrayDecoder::<$t, _>::new(
1054 ctx, $data_type, $tz,
1055 )))
1056 }};
1057 }
1058 macro_rules! decimal_decoder {
1059 ($t:ty, $p:expr, $s:expr) => {
1060 Ok(Box::new(DecimalArrayDecoder::<$t>::new(ctx, $p, $s)))
1061 };
1062 }
1063
1064 downcast_integer! {
1065 *data_type => (primitive_decoder, data_type),
1066 DataType::Null => Ok(Box::new(NullArrayDecoder::new(ctx))),
1067 DataType::Float16 => primitive_decoder!(Float16Type, data_type),
1068 DataType::Float32 => primitive_decoder!(Float32Type, data_type),
1069 DataType::Float64 => primitive_decoder!(Float64Type, data_type),
1070 DataType::Timestamp(TimeUnit::Second, None) => {
1071 timestamp_decoder!(TimestampSecondType, data_type, Utc)
1072 },
1073 DataType::Timestamp(TimeUnit::Millisecond, None) => {
1074 timestamp_decoder!(TimestampMillisecondType, data_type, Utc)
1075 },
1076 DataType::Timestamp(TimeUnit::Microsecond, None) => {
1077 timestamp_decoder!(TimestampMicrosecondType, data_type, Utc)
1078 },
1079 DataType::Timestamp(TimeUnit::Nanosecond, None) => {
1080 timestamp_decoder!(TimestampNanosecondType, data_type, Utc)
1081 },
1082 DataType::Timestamp(TimeUnit::Second, Some(ref tz)) => {
1083 let tz: Tz = tz.parse()?;
1084 timestamp_decoder!(TimestampSecondType, data_type, tz)
1085 },
1086 DataType::Timestamp(TimeUnit::Millisecond, Some(ref tz)) => {
1087 let tz: Tz = tz.parse()?;
1088 timestamp_decoder!(TimestampMillisecondType, data_type, tz)
1089 },
1090 DataType::Timestamp(TimeUnit::Microsecond, Some(ref tz)) => {
1091 let tz: Tz = tz.parse()?;
1092 timestamp_decoder!(TimestampMicrosecondType, data_type, tz)
1093 },
1094 DataType::Timestamp(TimeUnit::Nanosecond, Some(ref tz)) => {
1095 let tz: Tz = tz.parse()?;
1096 timestamp_decoder!(TimestampNanosecondType, data_type, tz)
1097 },
1098 DataType::Date32 => primitive_decoder!(Date32Type, data_type),
1099 DataType::Date64 => primitive_decoder!(Date64Type, data_type),
1100 DataType::Time32(TimeUnit::Second) => primitive_decoder!(Time32SecondType, data_type),
1101 DataType::Time32(TimeUnit::Millisecond) => primitive_decoder!(Time32MillisecondType, data_type),
1102 DataType::Time64(TimeUnit::Microsecond) => primitive_decoder!(Time64MicrosecondType, data_type),
1103 DataType::Time64(TimeUnit::Nanosecond) => primitive_decoder!(Time64NanosecondType, data_type),
1104 DataType::Duration(TimeUnit::Nanosecond) => primitive_decoder!(DurationNanosecondType, data_type),
1105 DataType::Duration(TimeUnit::Microsecond) => primitive_decoder!(DurationMicrosecondType, data_type),
1106 DataType::Duration(TimeUnit::Millisecond) => primitive_decoder!(DurationMillisecondType, data_type),
1107 DataType::Duration(TimeUnit::Second) => primitive_decoder!(DurationSecondType, data_type),
1108 DataType::Decimal32(p, s) => decimal_decoder!(Decimal32Type, p, s),
1109 DataType::Decimal64(p, s) => decimal_decoder!(Decimal64Type, p, s),
1110 DataType::Decimal128(p, s) => decimal_decoder!(Decimal128Type, p, s),
1111 DataType::Decimal256(p, s) => decimal_decoder!(Decimal256Type, p, s),
1112 DataType::Boolean => Ok(Box::new(BooleanArrayDecoder::new(ctx))),
1113 DataType::Utf8 => Ok(Box::new(StringArrayDecoder::<i32>::new(ctx))),
1114 DataType::Utf8View => Ok(Box::new(StringViewArrayDecoder::new(ctx))),
1115 DataType::LargeUtf8 => Ok(Box::new(StringArrayDecoder::<i64>::new(ctx))),
1116 DataType::List(_) => Ok(Box::new(ListArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
1117 DataType::LargeList(_) => Ok(Box::new(ListArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
1118 DataType::ListView(_) => Ok(Box::new(ListViewArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
1119 DataType::LargeListView(_) => Ok(Box::new(ListViewArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
1120 DataType::FixedSizeList(_, _) => Ok(Box::new(FixedSizeListArrayDecoder::new(ctx, data_type, is_nullable)?)),
1121 DataType::Struct(_) => Ok(Box::new(StructArrayDecoder::new(ctx, data_type, is_nullable)?)),
1122 DataType::Binary => Ok(Box::new(BinaryArrayDecoder::<i32>::default())),
1123 DataType::LargeBinary => Ok(Box::new(BinaryArrayDecoder::<i64>::default())),
1124 DataType::FixedSizeBinary(len) => Ok(Box::new(FixedSizeBinaryArrayDecoder::new(len))),
1125 DataType::BinaryView => Ok(Box::new(BinaryViewDecoder::default())),
1126 DataType::Map(_, _) => Ok(Box::new(MapArrayDecoder::new(ctx, data_type, is_nullable)?)),
1127 DataType::RunEndEncoded(ref r, _) => match r.data_type() {
1128 DataType::Int16 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int16Type>::new(ctx, data_type, is_nullable)?)),
1129 DataType::Int32 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int32Type>::new(ctx, data_type, is_nullable)?)),
1130 DataType::Int64 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int64Type>::new(ctx, data_type, is_nullable)?)),
1131 d => unreachable!("unsupported run end index type: {d}"),
1132 },
1133 _ => Err(ArrowError::NotYetImplemented(format!("Support for {data_type} in JSON reader")))
1134 }
1135}
1136
1137#[cfg(test)]
1138mod tests {
1139 use arrow_array::cast::AsArray;
1140 use arrow_array::{
1141 Array, BooleanArray, Float64Array, GenericListViewArray, Int32Array, ListArray, MapArray,
1142 NullArray, OffsetSizeTrait, StringArray, StringViewArray, StructArray,
1143 };
1144 use arrow_buffer::{ArrowNativeType, NullBuffer, OffsetBuffer, ScalarBuffer};
1145 use arrow_cast::display::{ArrayFormatter, FormatOptions};
1146 use arrow_schema::{Field, Fields};
1147 use serde_json::json;
1148 use std::fs::File;
1149 use std::io::{BufReader, Cursor, Seek};
1150 use std::sync::Mutex;
1151
1152 use super::*;
1153
1154 fn do_read(
1155 buf: &str,
1156 batch_size: usize,
1157 coerce_primitive: bool,
1158 strict_mode: bool,
1159 schema: SchemaRef,
1160 ) -> Vec<RecordBatch> {
1161 let config = ReaderBuilder::new(schema)
1162 .with_batch_size(batch_size)
1163 .with_strict_mode(strict_mode)
1164 .with_coerce_primitive(coerce_primitive);
1165 do_read_config(buf, config)
1166 }
1167
1168 fn do_read_config(buf: &str, builder: ReaderBuilder) -> Vec<RecordBatch> {
1169 let mut unbuffered = vec![];
1170
1171 for batch_size in [1, 3, 100, builder.batch_size] {
1173 unbuffered = builder
1174 .clone()
1175 .with_batch_size(batch_size)
1176 .build(Cursor::new(buf.as_bytes()))
1177 .unwrap()
1178 .collect::<Result<Vec<_>, _>>()
1179 .unwrap();
1180
1181 for b in unbuffered.iter().take(unbuffered.len() - 1) {
1182 assert_eq!(b.num_rows(), batch_size)
1183 }
1184
1185 for b in [1, 3, 5] {
1187 let buffered = builder
1188 .clone()
1189 .with_batch_size(batch_size)
1190 .build(BufReader::with_capacity(b, Cursor::new(buf.as_bytes())))
1191 .unwrap()
1192 .collect::<Result<Vec<_>, _>>()
1193 .unwrap();
1194 assert_eq!(unbuffered, buffered);
1195 }
1196 }
1197
1198 unbuffered
1199 }
1200
1201 #[test]
1202 fn test_basic() {
1203 let buf = r#"
1204 {"a": 1, "b": 2, "c": true, "d": 1}
1205 {"a": 2E0, "b": 4, "c": false, "d": 2, "e": 254}
1206
1207 {"b": 6, "a": 2.0, "d": 45}
1208 {"b": "5", "a": 2}
1209 {"b": 4e0}
1210 {"b": 7, "a": null}
1211 "#;
1212
1213 let schema = Arc::new(Schema::new(vec![
1214 Field::new("a", DataType::Int64, true),
1215 Field::new("b", DataType::Int32, true),
1216 Field::new("c", DataType::Boolean, true),
1217 Field::new("d", DataType::Date32, true),
1218 Field::new("e", DataType::Date64, true),
1219 ]));
1220
1221 let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
1222 assert!(decoder.is_empty());
1223 assert_eq!(decoder.len(), 0);
1224 assert!(!decoder.has_partial_record());
1225 assert_eq!(decoder.decode(buf.as_bytes()).unwrap(), 221);
1226 assert!(!decoder.is_empty());
1227 assert_eq!(decoder.len(), 6);
1228 assert!(!decoder.has_partial_record());
1229 let batch = decoder.flush().unwrap().unwrap();
1230 assert_eq!(batch.num_rows(), 6);
1231 assert!(decoder.is_empty());
1232 assert_eq!(decoder.len(), 0);
1233 assert!(!decoder.has_partial_record());
1234
1235 let batches = do_read(buf, 1024, false, false, schema);
1236 assert_eq!(batches.len(), 1);
1237
1238 let col1 = batches[0].column(0).as_primitive::<Int64Type>();
1239 assert_eq!(col1.null_count(), 2);
1240 assert_eq!(col1.values(), &[1, 2, 2, 2, 0, 0]);
1241 assert!(col1.is_null(4));
1242 assert!(col1.is_null(5));
1243
1244 let col2 = batches[0].column(1).as_primitive::<Int32Type>();
1245 assert_eq!(col2.null_count(), 0);
1246 assert_eq!(col2.values(), &[2, 4, 6, 5, 4, 7]);
1247
1248 let col3 = batches[0].column(2).as_boolean();
1249 assert_eq!(col3.null_count(), 4);
1250 assert!(col3.value(0));
1251 assert!(!col3.is_null(0));
1252 assert!(!col3.value(1));
1253 assert!(!col3.is_null(1));
1254
1255 let col4 = batches[0].column(3).as_primitive::<Date32Type>();
1256 assert_eq!(col4.null_count(), 3);
1257 assert!(col4.is_null(3));
1258 assert_eq!(col4.values(), &[1, 2, 45, 0, 0, 0]);
1259
1260 let col5 = batches[0].column(4).as_primitive::<Date64Type>();
1261 assert_eq!(col5.null_count(), 5);
1262 assert!(col5.is_null(0));
1263 assert!(col5.is_null(2));
1264 assert!(col5.is_null(3));
1265 assert_eq!(col5.values(), &[0, 254, 0, 0, 0, 0]);
1266 }
1267
1268 #[test]
1269 fn test_string() {
1270 let buf = r#"
1271 {"a": "1", "b": "2"}
1272 {"a": "hello", "b": "shoo"}
1273 {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
1274
1275 {"b": null}
1276 {"b": "", "a": null}
1277
1278 "#;
1279 let schema = Arc::new(Schema::new(vec![
1280 Field::new("a", DataType::Utf8, true),
1281 Field::new("b", DataType::LargeUtf8, true),
1282 ]));
1283
1284 let batches = do_read(buf, 1024, false, false, schema);
1285 assert_eq!(batches.len(), 1);
1286
1287 let col1 = batches[0].column(0).as_string::<i32>();
1288 assert_eq!(col1.null_count(), 2);
1289 assert_eq!(col1.value(0), "1");
1290 assert_eq!(col1.value(1), "hello");
1291 assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1292 assert!(col1.is_null(3));
1293 assert!(col1.is_null(4));
1294
1295 let col2 = batches[0].column(1).as_string::<i64>();
1296 assert_eq!(col2.null_count(), 1);
1297 assert_eq!(col2.value(0), "2");
1298 assert_eq!(col2.value(1), "shoo");
1299 assert_eq!(col2.value(2), "\t😁foo");
1300 assert!(col2.is_null(3));
1301 assert_eq!(col2.value(4), "");
1302 }
1303
1304 #[test]
1305 fn test_long_string_view_allocation() {
1306 let expected_capacity: usize = 41;
1316
1317 let buf = r#"
1318 {"a": "short", "b": "dummy"}
1319 {"a": "this is definitely long", "b": "dummy"}
1320 {"a": "hello", "b": "dummy"}
1321 {"a": "\nfoobar😀asfgÿ", "b": "dummy"}
1322 "#;
1323
1324 let schema = Arc::new(Schema::new(vec![
1325 Field::new("a", DataType::Utf8View, true),
1326 Field::new("b", DataType::LargeUtf8, true),
1327 ]));
1328
1329 let batches = do_read(buf, 1024, false, false, schema);
1330 assert_eq!(batches.len(), 1, "Expected one record batch");
1331
1332 let col_a = batches[0].column(0);
1334 let string_view_array = col_a
1335 .as_any()
1336 .downcast_ref::<StringViewArray>()
1337 .expect("Column should be a StringViewArray");
1338
1339 let data_buffer = string_view_array.to_data().buffers()[0].len();
1342
1343 assert!(
1346 data_buffer >= expected_capacity,
1347 "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1348 );
1349
1350 assert_eq!(string_view_array.value(0), "short");
1352 assert_eq!(string_view_array.value(1), "this is definitely long");
1353 assert_eq!(string_view_array.value(2), "hello");
1354 assert_eq!(string_view_array.value(3), "\nfoobar😀asfgÿ");
1355 }
1356
1357 #[test]
1359 fn test_numeric_view_allocation() {
1360 let expected_capacity: usize = 33;
1368
1369 let buf = r#"
1370 {"n": 123456789}
1371 {"n": 1000000000000}
1372 {"n": 3.1415}
1373 {"n": 2.718281828459045}
1374 "#;
1375
1376 let schema = Arc::new(Schema::new(vec![Field::new("n", DataType::Utf8View, true)]));
1377
1378 let batches = do_read(buf, 1024, true, false, schema);
1379 assert_eq!(batches.len(), 1, "Expected one record batch");
1380
1381 let col_n = batches[0].column(0);
1382 let string_view_array = col_n
1383 .as_any()
1384 .downcast_ref::<StringViewArray>()
1385 .expect("Column should be a StringViewArray");
1386
1387 let data_buffer = string_view_array.to_data().buffers()[0].len();
1389 assert!(
1390 data_buffer >= expected_capacity,
1391 "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1392 );
1393
1394 assert_eq!(string_view_array.value(0), "123456789");
1397 assert_eq!(string_view_array.value(1), "1000000000000");
1398 assert_eq!(string_view_array.value(2), "3.1415");
1399 assert_eq!(string_view_array.value(3), "2.718281828459045");
1400 }
1401
1402 #[test]
1403 fn test_string_with_uft8view() {
1404 let buf = r#"
1405 {"a": "1", "b": "2"}
1406 {"a": "hello", "b": "shoo"}
1407 {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
1408
1409 {"b": null}
1410 {"b": "", "a": null}
1411
1412 "#;
1413 let schema = Arc::new(Schema::new(vec![
1414 Field::new("a", DataType::Utf8View, true),
1415 Field::new("b", DataType::LargeUtf8, true),
1416 ]));
1417
1418 let batches = do_read(buf, 1024, false, false, schema);
1419 assert_eq!(batches.len(), 1);
1420
1421 let col1 = batches[0].column(0).as_string_view();
1422 assert_eq!(col1.null_count(), 2);
1423 assert_eq!(col1.value(0), "1");
1424 assert_eq!(col1.value(1), "hello");
1425 assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1426 assert!(col1.is_null(3));
1427 assert!(col1.is_null(4));
1428 assert_eq!(col1.data_type(), &DataType::Utf8View);
1429
1430 let col2 = batches[0].column(1).as_string::<i64>();
1431 assert_eq!(col2.null_count(), 1);
1432 assert_eq!(col2.value(0), "2");
1433 assert_eq!(col2.value(1), "shoo");
1434 assert_eq!(col2.value(2), "\t😁foo");
1435 assert!(col2.is_null(3));
1436 assert_eq!(col2.value(4), "");
1437 }
1438
1439 #[test]
1440 fn test_complex() {
1441 let buf = r#"
1442 {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3}, {"c": 4}]}}
1443 {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1444 {"list": null, "nested": {"a": null}}
1445 "#;
1446
1447 let schema = Arc::new(Schema::new(vec![
1448 Field::new_list("list", Field::new("element", DataType::Int32, false), true),
1449 Field::new_struct(
1450 "nested",
1451 vec![
1452 Field::new("a", DataType::Int32, true),
1453 Field::new("b", DataType::Int32, true),
1454 ],
1455 true,
1456 ),
1457 Field::new_struct(
1458 "nested_list",
1459 vec![Field::new_list(
1460 "list2",
1461 Field::new_struct(
1462 "element",
1463 vec![Field::new("c", DataType::Int32, false)],
1464 false,
1465 ),
1466 true,
1467 )],
1468 true,
1469 ),
1470 ]));
1471
1472 let batches = do_read(buf, 1024, false, false, schema);
1473 assert_eq!(batches.len(), 1);
1474
1475 let list = batches[0].column(0).as_list::<i32>();
1476 assert_eq!(list.len(), 3);
1477 assert_eq!(list.value_offsets(), &[0, 0, 2, 2]);
1478 assert_eq!(list.null_count(), 1);
1479 assert!(list.is_null(2));
1480 let list_values = list.values().as_primitive::<Int32Type>();
1481 assert_eq!(list_values.values(), &[5, 6]);
1482
1483 let nested = batches[0].column(1).as_struct();
1484 let a = nested.column(0).as_primitive::<Int32Type>();
1485 assert_eq!(list.null_count(), 1);
1486 assert_eq!(a.values(), &[1, 7, 0]);
1487 assert!(list.is_null(2));
1488
1489 let b = nested.column(1).as_primitive::<Int32Type>();
1490 assert_eq!(b.null_count(), 2);
1491 assert_eq!(b.len(), 3);
1492 assert_eq!(b.value(0), 2);
1493 assert!(b.is_null(1));
1494 assert!(b.is_null(2));
1495
1496 let nested_list = batches[0].column(2).as_struct();
1497 assert_eq!(nested_list.len(), 3);
1498 assert_eq!(nested_list.null_count(), 1);
1499 assert!(nested_list.is_null(2));
1500
1501 let list2 = nested_list.column(0).as_list::<i32>();
1502 assert_eq!(list2.len(), 3);
1503 assert_eq!(list2.null_count(), 1);
1504 assert_eq!(list2.value_offsets(), &[0, 2, 2, 2]);
1505 assert!(list2.is_null(2));
1506
1507 let list2_values = list2.values().as_struct();
1508
1509 let c = list2_values.column(0).as_primitive::<Int32Type>();
1510 assert_eq!(c.values(), &[3, 4]);
1511 }
1512
1513 #[test]
1514 fn test_projection() {
1515 let buf = r#"
1516 {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3, "d": 5}, {"c": 4}]}}
1517 {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1518 "#;
1519
1520 let schema = Arc::new(Schema::new(vec![
1521 Field::new_struct(
1522 "nested",
1523 vec![Field::new("a", DataType::Int32, false)],
1524 true,
1525 ),
1526 Field::new_struct(
1527 "nested_list",
1528 vec![Field::new_list(
1529 "list2",
1530 Field::new_struct(
1531 "element",
1532 vec![Field::new("d", DataType::Int32, true)],
1533 false,
1534 ),
1535 true,
1536 )],
1537 true,
1538 ),
1539 ]));
1540
1541 let batches = do_read(buf, 1024, false, false, schema);
1542 assert_eq!(batches.len(), 1);
1543
1544 let nested = batches[0].column(0).as_struct();
1545 assert_eq!(nested.num_columns(), 1);
1546 let a = nested.column(0).as_primitive::<Int32Type>();
1547 assert_eq!(a.null_count(), 0);
1548 assert_eq!(a.values(), &[1, 7]);
1549
1550 let nested_list = batches[0].column(1).as_struct();
1551 assert_eq!(nested_list.num_columns(), 1);
1552 assert_eq!(nested_list.null_count(), 0);
1553
1554 let list2 = nested_list.column(0).as_list::<i32>();
1555 assert_eq!(list2.value_offsets(), &[0, 2, 2]);
1556 assert_eq!(list2.null_count(), 0);
1557
1558 let child = list2.values().as_struct();
1559 assert_eq!(child.num_columns(), 1);
1560 assert_eq!(child.len(), 2);
1561 assert_eq!(child.null_count(), 0);
1562
1563 let c = child.column(0).as_primitive::<Int32Type>();
1564 assert_eq!(c.values(), &[5, 0]);
1565 assert_eq!(c.null_count(), 1);
1566 assert!(c.is_null(1));
1567 }
1568
1569 #[test]
1570 fn test_map() {
1571 let buf = r#"
1572 {"map": {"a": ["foo", null]}}
1573 {"map": {"a": [null], "b": []}}
1574 {"map": {"c": null, "a": ["baz"]}}
1575 "#;
1576 let map = Field::new_map(
1577 "map",
1578 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1579 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1580 Field::new_list(
1581 Field::MAP_VALUE_FIELD_DEFAULT_NAME,
1582 Field::new("element", DataType::Utf8, true),
1583 true,
1584 ),
1585 false,
1586 true,
1587 );
1588
1589 let schema = Arc::new(Schema::new(vec![map]));
1590
1591 let batches = do_read(buf, 1024, false, false, schema);
1592 assert_eq!(batches.len(), 1);
1593
1594 let map = batches[0].column(0).as_map();
1595 let map_keys = map.keys().as_string::<i32>();
1596 let map_values = map.values().as_list::<i32>();
1597 assert_eq!(map.value_offsets(), &[0, 1, 3, 5]);
1598
1599 let k: Vec<_> = map_keys.iter().flatten().collect();
1600 assert_eq!(&k, &["a", "a", "b", "c", "a"]);
1601
1602 let list_values = map_values.values().as_string::<i32>();
1603 let lv: Vec<_> = list_values.iter().collect();
1604 assert_eq!(&lv, &[Some("foo"), None, None, Some("baz")]);
1605 assert_eq!(map_values.value_offsets(), &[0, 2, 3, 3, 3, 4]);
1606 assert_eq!(map_values.null_count(), 1);
1607 assert!(map_values.is_null(3));
1608
1609 let options = FormatOptions::default().with_null("null");
1610 let formatter = ArrayFormatter::try_new(map, &options).unwrap();
1611 assert_eq!(formatter.value(0).to_string(), "{a: [foo, null]}");
1612 assert_eq!(formatter.value(1).to_string(), "{a: [null], b: []}");
1613 assert_eq!(formatter.value(2).to_string(), "{c: null, a: [baz]}");
1614 }
1615
1616 #[test]
1617 fn test_map_non_nullable_value() {
1618 let map = Field::new_map(
1619 "map",
1620 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1621 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1622 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1623 false,
1624 false,
1625 );
1626 let schema = Arc::new(Schema::new(vec![map]));
1627 let buf = r#"{"map": {"key": null}}"#;
1628
1629 let err = ReaderBuilder::new(schema)
1630 .build(Cursor::new(buf.as_bytes()))
1631 .unwrap()
1632 .read()
1633 .unwrap_err();
1634
1635 assert_eq!(
1636 err.to_string(),
1637 "Invalid argument error: Found unmasked nulls for non-nullable StructArray field \"value\""
1638 );
1639 }
1640
1641 #[test]
1642 fn test_not_coercing_primitive_into_string_without_flag() {
1643 let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, true)]));
1644
1645 let buf = r#"{"a": 1}"#;
1646 let err = ReaderBuilder::new(schema.clone())
1647 .with_batch_size(1024)
1648 .build(Cursor::new(buf.as_bytes()))
1649 .unwrap()
1650 .read()
1651 .unwrap_err();
1652
1653 assert_eq!(
1654 err.to_string(),
1655 "Json error: whilst decoding field 'a': expected string got 1"
1656 );
1657
1658 let buf = r#"{"a": true}"#;
1659 let err = ReaderBuilder::new(schema)
1660 .with_batch_size(1024)
1661 .build(Cursor::new(buf.as_bytes()))
1662 .unwrap()
1663 .read()
1664 .unwrap_err();
1665
1666 assert_eq!(
1667 err.to_string(),
1668 "Json error: whilst decoding field 'a': expected string got true"
1669 );
1670 }
1671
1672 #[test]
1673 fn test_coercing_primitive_into_string() {
1674 let buf = r#"
1675 {"a": 1, "b": 2, "c": true}
1676 {"a": 2E0, "b": 4, "c": false}
1677
1678 {"b": 6, "a": 2.0}
1679 {"b": "5", "a": 2}
1680 {"b": 4e0}
1681 {"b": 7, "a": null}
1682 "#;
1683
1684 let schema = Arc::new(Schema::new(vec![
1685 Field::new("a", DataType::Utf8, true),
1686 Field::new("b", DataType::Utf8, true),
1687 Field::new("c", DataType::Utf8, true),
1688 ]));
1689
1690 let batches = do_read(buf, 1024, true, false, schema);
1691 assert_eq!(batches.len(), 1);
1692
1693 let col1 = batches[0].column(0).as_string::<i32>();
1694 assert_eq!(col1.null_count(), 2);
1695 assert_eq!(col1.value(0), "1");
1696 assert_eq!(col1.value(1), "2E0");
1697 assert_eq!(col1.value(2), "2.0");
1698 assert_eq!(col1.value(3), "2");
1699 assert!(col1.is_null(4));
1700 assert!(col1.is_null(5));
1701
1702 let col2 = batches[0].column(1).as_string::<i32>();
1703 assert_eq!(col2.null_count(), 0);
1704 assert_eq!(col2.value(0), "2");
1705 assert_eq!(col2.value(1), "4");
1706 assert_eq!(col2.value(2), "6");
1707 assert_eq!(col2.value(3), "5");
1708 assert_eq!(col2.value(4), "4e0");
1709 assert_eq!(col2.value(5), "7");
1710
1711 let col3 = batches[0].column(2).as_string::<i32>();
1712 assert_eq!(col3.null_count(), 4);
1713 assert_eq!(col3.value(0), "true");
1714 assert_eq!(col3.value(1), "false");
1715 assert!(col3.is_null(2));
1716 assert!(col3.is_null(3));
1717 assert!(col3.is_null(4));
1718 assert!(col3.is_null(5));
1719 }
1720
1721 fn test_decimal<T: DecimalType>(data_type: DataType) {
1722 let buf = r#"
1723 {"a": 1, "b": 2, "c": 38.30}
1724 {"a": 2, "b": 4, "c": 123.456}
1725
1726 {"b": 1337, "a": "2.0452"}
1727 {"b": "5", "a": "11034.2"}
1728 {"b": 40}
1729 {"b": 1234, "a": null}
1730 "#;
1731
1732 let schema = Arc::new(Schema::new(vec![
1733 Field::new("a", data_type.clone(), true),
1734 Field::new("b", data_type.clone(), true),
1735 Field::new("c", data_type, true),
1736 ]));
1737
1738 let batches = do_read(buf, 1024, true, false, schema);
1739 assert_eq!(batches.len(), 1);
1740
1741 let col1 = batches[0].column(0).as_primitive::<T>();
1742 assert_eq!(col1.null_count(), 2);
1743 assert!(col1.is_null(4));
1744 assert!(col1.is_null(5));
1745 assert_eq!(
1746 col1.values(),
1747 &[100, 200, 205, 1103420, 0, 0].map(T::Native::usize_as)
1748 );
1749
1750 let col2 = batches[0].column(1).as_primitive::<T>();
1751 assert_eq!(col2.null_count(), 0);
1752 assert_eq!(
1753 col2.values(),
1754 &[200, 400, 133700, 500, 4000, 123400].map(T::Native::usize_as)
1755 );
1756
1757 let col3 = batches[0].column(2).as_primitive::<T>();
1758 assert_eq!(col3.null_count(), 4);
1759 assert!(!col3.is_null(0));
1760 assert!(!col3.is_null(1));
1761 assert!(col3.is_null(2));
1762 assert!(col3.is_null(3));
1763 assert!(col3.is_null(4));
1764 assert!(col3.is_null(5));
1765 assert_eq!(
1766 col3.values(),
1767 &[3830, 12346, 0, 0, 0, 0].map(T::Native::usize_as)
1768 );
1769 }
1770
1771 #[test]
1772 fn test_decimal_number_formats() {
1773 let buf = r#"
1777 {"a": 0e0, "b": " 1.5 ", "c": 1234.5}
1778 {"a": 1.5E2, "b": "-0.005", "c": -150}
1779 {"a": 1e-3, "b": "1e2", "c": 1E2}
1780 {"a": -0E+0, "b": "+.5", "c": 5}
1781 "#;
1782 let schema = Arc::new(Schema::new(vec![
1783 Field::new("a", DataType::Decimal128(10, 2), true),
1784 Field::new("b", DataType::Decimal64(18, 2), true),
1785 Field::new("c", DataType::Decimal128(10, -2), true),
1786 ]));
1787 let batches = do_read(buf, 1024, true, false, schema);
1788 assert_eq!(batches.len(), 1);
1789 assert_eq!(
1790 batches[0]
1791 .column(0)
1792 .as_primitive::<Decimal128Type>()
1793 .values(),
1794 &[0, 15000, 0, 0]
1795 );
1796 assert_eq!(
1797 batches[0]
1798 .column(1)
1799 .as_primitive::<Decimal64Type>()
1800 .values(),
1801 &[150, -1, 10000, 50]
1802 );
1803 assert_eq!(
1804 batches[0]
1805 .column(2)
1806 .as_primitive::<Decimal128Type>()
1807 .values(),
1808 &[12, -2, 1, 0]
1809 );
1810
1811 let schema = Arc::new(Schema::new(vec![Field::new(
1814 "a",
1815 DataType::Decimal128(5, 2),
1816 true,
1817 )]));
1818 let long = "1".repeat(300);
1819 for (buf, expected) in [
1820 (
1821 r#"{"a": "abc"}"#.to_string(),
1822 "Invalid decimal format: \"abc\"",
1823 ),
1824 (
1825 r#"{"a": 123456789}"#.to_string(),
1826 "does not fit in Decimal128(5, 2)",
1827 ),
1828 (
1829 r#"{"a": 1e99999}"#.to_string(),
1830 "does not fit in Decimal128(5, 2)",
1831 ),
1832 (
1833 format!(r#"{{"a": {long}}}"#),
1834 "does not fit in Decimal128(5, 2)",
1835 ),
1836 (
1837 r#"{"a": 4825037936439135476.2609835314269495255615E-14}"#.to_string(),
1838 "does not fit in Decimal128(5, 2)",
1839 ),
1840 ] {
1841 let err = ReaderBuilder::new(schema.clone())
1842 .build(Cursor::new(buf.as_bytes()))
1843 .unwrap()
1844 .next()
1845 .unwrap()
1846 .unwrap_err()
1847 .to_string();
1848 assert!(err.contains(expected), "{buf}: {err}");
1849
1850 let batch = ReaderBuilder::new(schema.clone())
1851 .with_ignore_type_conflicts(true)
1852 .build(Cursor::new(buf.as_bytes()))
1853 .unwrap()
1854 .next()
1855 .unwrap()
1856 .unwrap();
1857 assert!(batch.column(0).is_null(0), "{buf}");
1858 }
1859 }
1860
1861 #[test]
1862 fn test_decimals() {
1863 test_decimal::<Decimal32Type>(DataType::Decimal32(8, 2));
1864 test_decimal::<Decimal64Type>(DataType::Decimal64(10, 2));
1865 test_decimal::<Decimal128Type>(DataType::Decimal128(10, 2));
1866 test_decimal::<Decimal256Type>(DataType::Decimal256(10, 2));
1867 }
1868
1869 fn test_timestamp<T: ArrowTimestampType>() {
1870 let buf = r#"
1871 {"a": 1, "b": "2020-09-08T13:42:29.190855+00:00", "c": 38.30, "d": "1997-01-31T09:26:56.123"}
1872 {"a": 2, "b": "2020-09-08T13:42:29.190855Z", "c": 123.456, "d": 123.456}
1873
1874 {"b": 1337, "b": "2020-09-08T13:42:29Z", "c": "1997-01-31T09:26:56.123", "d": "1997-01-31T09:26:56.123Z"}
1875 {"b": 40, "c": "2020-09-08T13:42:29.190855+00:00", "d": "1997-01-31 09:26:56.123-05:00"}
1876 {"b": 1234, "a": null, "c": "1997-01-31 09:26:56.123Z", "d": "1997-01-31 092656"}
1877 {"c": "1997-01-31T14:26:56.123-05:00", "d": "1997-01-31"}
1878 "#;
1879
1880 let with_timezone = DataType::Timestamp(T::UNIT, Some("+08:00".into()));
1881 let schema = Arc::new(Schema::new(vec![
1882 Field::new("a", T::DATA_TYPE, true),
1883 Field::new("b", T::DATA_TYPE, true),
1884 Field::new("c", T::DATA_TYPE, true),
1885 Field::new("d", with_timezone, true),
1886 ]));
1887
1888 let batches = do_read(buf, 1024, true, false, schema);
1889 assert_eq!(batches.len(), 1);
1890
1891 let unit_in_nanos: i64 = match T::UNIT {
1892 TimeUnit::Second => 1_000_000_000,
1893 TimeUnit::Millisecond => 1_000_000,
1894 TimeUnit::Microsecond => 1_000,
1895 TimeUnit::Nanosecond => 1,
1896 };
1897
1898 let col1 = batches[0].column(0).as_primitive::<T>();
1899 assert_eq!(col1.null_count(), 4);
1900 assert!(col1.is_null(2));
1901 assert!(col1.is_null(3));
1902 assert!(col1.is_null(4));
1903 assert!(col1.is_null(5));
1904 assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1905
1906 let col2 = batches[0].column(1).as_primitive::<T>();
1907 assert_eq!(col2.null_count(), 1);
1908 assert!(col2.is_null(5));
1909 assert_eq!(
1910 col2.values(),
1911 &[
1912 1599572549190855000 / unit_in_nanos,
1913 1599572549190855000 / unit_in_nanos,
1914 1599572549000000000 / unit_in_nanos,
1915 40,
1916 1234,
1917 0
1918 ]
1919 );
1920
1921 let col3 = batches[0].column(2).as_primitive::<T>();
1922 assert_eq!(col3.null_count(), 0);
1923 assert_eq!(
1924 col3.values(),
1925 &[
1926 38,
1927 123,
1928 854702816123000000 / unit_in_nanos,
1929 1599572549190855000 / unit_in_nanos,
1930 854702816123000000 / unit_in_nanos,
1931 854738816123000000 / unit_in_nanos
1932 ]
1933 );
1934
1935 let col4 = batches[0].column(3).as_primitive::<T>();
1936
1937 assert_eq!(col4.null_count(), 0);
1938 assert_eq!(
1939 col4.values(),
1940 &[
1941 854674016123000000 / unit_in_nanos,
1942 123,
1943 854702816123000000 / unit_in_nanos,
1944 854720816123000000 / unit_in_nanos,
1945 854674016000000000 / unit_in_nanos,
1946 854640000000000000 / unit_in_nanos
1947 ]
1948 );
1949 }
1950
1951 #[test]
1952 #[cfg_attr(miri, ignore)] fn test_timestamps() {
1954 test_timestamp::<TimestampSecondType>();
1955 test_timestamp::<TimestampMillisecondType>();
1956 test_timestamp::<TimestampMicrosecondType>();
1957 test_timestamp::<TimestampNanosecondType>();
1958 }
1959
1960 fn test_time<T: ArrowTemporalType>() {
1961 let buf = r#"
1962 {"a": 1, "b": "09:26:56.123 AM", "c": 38.30}
1963 {"a": 2, "b": "23:59:59", "c": 123.456}
1964
1965 {"b": 1337, "b": "6:00 pm", "c": "09:26:56.123"}
1966 {"b": 40, "c": "13:42:29.190855"}
1967 {"b": 1234, "a": null, "c": "09:26:56.123"}
1968 {"c": "14:26:56.123"}
1969 "#;
1970
1971 let (DataType::Time32(unit) | DataType::Time64(unit)) = T::DATA_TYPE else {
1972 unreachable!()
1973 };
1974
1975 let unit_in_nanos = match unit {
1976 TimeUnit::Second => 1_000_000_000,
1977 TimeUnit::Millisecond => 1_000_000,
1978 TimeUnit::Microsecond => 1_000,
1979 TimeUnit::Nanosecond => 1,
1980 };
1981
1982 let schema = Arc::new(Schema::new(vec![
1983 Field::new("a", T::DATA_TYPE, true),
1984 Field::new("b", T::DATA_TYPE, true),
1985 Field::new("c", T::DATA_TYPE, true),
1986 ]));
1987
1988 let batches = do_read(buf, 1024, true, false, schema);
1989 assert_eq!(batches.len(), 1);
1990
1991 let col1 = batches[0].column(0).as_primitive::<T>();
1992 assert_eq!(col1.null_count(), 4);
1993 assert!(col1.is_null(2));
1994 assert!(col1.is_null(3));
1995 assert!(col1.is_null(4));
1996 assert!(col1.is_null(5));
1997 assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1998
1999 let col2 = batches[0].column(1).as_primitive::<T>();
2000 assert_eq!(col2.null_count(), 1);
2001 assert!(col2.is_null(5));
2002 assert_eq!(
2003 col2.values(),
2004 &[
2005 34016123000000 / unit_in_nanos,
2006 86399000000000 / unit_in_nanos,
2007 64800000000000 / unit_in_nanos,
2008 40,
2009 1234,
2010 0
2011 ]
2012 .map(T::Native::usize_as)
2013 );
2014
2015 let col3 = batches[0].column(2).as_primitive::<T>();
2016 assert_eq!(col3.null_count(), 0);
2017 assert_eq!(
2018 col3.values(),
2019 &[
2020 38,
2021 123,
2022 34016123000000 / unit_in_nanos,
2023 49349190855000 / unit_in_nanos,
2024 34016123000000 / unit_in_nanos,
2025 52016123000000 / unit_in_nanos
2026 ]
2027 .map(T::Native::usize_as)
2028 );
2029 }
2030
2031 #[test]
2032 fn test_times() {
2033 test_time::<Time32MillisecondType>();
2034 test_time::<Time32SecondType>();
2035 test_time::<Time64MicrosecondType>();
2036 test_time::<Time64NanosecondType>();
2037 }
2038
2039 fn test_duration<T: ArrowTemporalType>() {
2040 let buf = r#"
2041 {"a": 1, "b": "2"}
2042 {"a": 3, "b": null}
2043 "#;
2044
2045 let schema = Arc::new(Schema::new(vec![
2046 Field::new("a", T::DATA_TYPE, true),
2047 Field::new("b", T::DATA_TYPE, true),
2048 ]));
2049
2050 let batches = do_read(buf, 1024, true, false, schema);
2051 assert_eq!(batches.len(), 1);
2052
2053 let col_a = batches[0].column_by_name("a").unwrap().as_primitive::<T>();
2054 assert_eq!(col_a.null_count(), 0);
2055 assert_eq!(col_a.values(), &[1, 3].map(T::Native::usize_as));
2056
2057 let col2 = batches[0].column_by_name("b").unwrap().as_primitive::<T>();
2058 assert_eq!(col2.null_count(), 1);
2059 assert_eq!(col2.values(), &[2, 0].map(T::Native::usize_as));
2060 }
2061
2062 #[test]
2063 fn test_durations() {
2064 test_duration::<DurationNanosecondType>();
2065 test_duration::<DurationMicrosecondType>();
2066 test_duration::<DurationMillisecondType>();
2067 test_duration::<DurationSecondType>();
2068 }
2069
2070 #[test]
2071 fn test_delta_checkpoint() {
2072 let json = "{\"protocol\":{\"minReaderVersion\":1,\"minWriterVersion\":2}}";
2073 let schema = Arc::new(Schema::new(vec![
2074 Field::new_struct(
2075 "protocol",
2076 vec![
2077 Field::new("minReaderVersion", DataType::Int32, true),
2078 Field::new("minWriterVersion", DataType::Int32, true),
2079 ],
2080 true,
2081 ),
2082 Field::new_struct(
2083 "add",
2084 vec![Field::new_map(
2085 "partitionValues",
2086 "key_value",
2087 Field::new("key", DataType::Utf8, false),
2088 Field::new("value", DataType::Utf8, true),
2089 false,
2090 false,
2091 )],
2092 true,
2093 ),
2094 ]));
2095
2096 let batches = do_read(json, 1024, true, false, schema);
2097 assert_eq!(batches.len(), 1);
2098
2099 let s: StructArray = batches.into_iter().next().unwrap().into();
2100 let opts = FormatOptions::default().with_null("null");
2101 let formatter = ArrayFormatter::try_new(&s, &opts).unwrap();
2102 assert_eq!(
2103 formatter.value(0).to_string(),
2104 "{protocol: {minReaderVersion: 1, minWriterVersion: 2}, add: null}"
2105 );
2106 }
2107
2108 #[test]
2109 fn struct_nullability() {
2110 let do_test = |child: DataType| {
2111 let non_null = r#"{"foo": {}}"#;
2113 let schema = Arc::new(Schema::new(vec![Field::new_struct(
2114 "foo",
2115 vec![Field::new("bar", child, false)],
2116 true,
2117 )]));
2118 let mut reader = ReaderBuilder::new(schema.clone())
2119 .build(Cursor::new(non_null.as_bytes()))
2120 .unwrap();
2121 assert!(reader.next().unwrap().is_err()); let null = r#"{"foo": {bar: null}}"#;
2124 let mut reader = ReaderBuilder::new(schema.clone())
2125 .build(Cursor::new(null.as_bytes()))
2126 .unwrap();
2127 assert!(reader.next().unwrap().is_err()); let null = r#"{"foo": null}"#;
2131 let mut reader = ReaderBuilder::new(schema)
2132 .build(Cursor::new(null.as_bytes()))
2133 .unwrap();
2134 let batch = reader.next().unwrap().unwrap();
2135 assert_eq!(batch.num_columns(), 1);
2136 let foo = batch.column(0).as_struct();
2137 assert_eq!(foo.len(), 1);
2138 assert!(foo.is_null(0));
2139 assert_eq!(foo.num_columns(), 1);
2140
2141 let bar = foo.column(0);
2142 assert_eq!(bar.len(), 1);
2143 assert!(bar.is_null(0));
2145 };
2146
2147 do_test(DataType::Boolean);
2148 do_test(DataType::Int32);
2149 do_test(DataType::Utf8);
2150 do_test(DataType::Decimal128(2, 1));
2151 do_test(DataType::Timestamp(
2152 TimeUnit::Microsecond,
2153 Some("+00:00".into()),
2154 ));
2155 }
2156
2157 #[test]
2158 fn test_truncation() {
2159 let buf = r#"
2160 {"i64": 9223372036854775807, "u64": 18446744073709551615 }
2161 {"i64": "9223372036854775807", "u64": "18446744073709551615" }
2162 {"i64": -9223372036854775808, "u64": 0 }
2163 {"i64": "-9223372036854775808", "u64": 0 }
2164 "#;
2165
2166 let schema = Arc::new(Schema::new(vec![
2167 Field::new("i64", DataType::Int64, true),
2168 Field::new("u64", DataType::UInt64, true),
2169 ]));
2170
2171 let batches = do_read(buf, 1024, true, false, schema);
2172 assert_eq!(batches.len(), 1);
2173
2174 let i64 = batches[0].column(0).as_primitive::<Int64Type>();
2175 assert_eq!(i64.values(), &[i64::MAX, i64::MAX, i64::MIN, i64::MIN]);
2176
2177 let u64 = batches[0].column(1).as_primitive::<UInt64Type>();
2178 assert_eq!(u64.values(), &[u64::MAX, u64::MAX, u64::MIN, u64::MIN]);
2179 }
2180
2181 #[test]
2182 fn test_timestamp_truncation() {
2183 let buf = r#"
2184 {"time": 9223372036854775807 }
2185 {"time": -9223372036854775808 }
2186 {"time": 9e5 }
2187 "#;
2188
2189 let schema = Arc::new(Schema::new(vec![Field::new(
2190 "time",
2191 DataType::Timestamp(TimeUnit::Nanosecond, None),
2192 true,
2193 )]));
2194
2195 let batches = do_read(buf, 1024, true, false, schema);
2196 assert_eq!(batches.len(), 1);
2197
2198 let i64 = batches[0]
2199 .column(0)
2200 .as_primitive::<TimestampNanosecondType>();
2201 assert_eq!(i64.values(), &[i64::MAX, i64::MIN, 900000]);
2202 }
2203
2204 #[test]
2205 fn test_strict_mode_no_missing_columns_in_schema() {
2206 let buf = r#"
2207 {"a": 1, "b": "2", "c": true}
2208 {"a": 2E0, "b": "4", "c": false}
2209 "#;
2210
2211 let schema = Arc::new(Schema::new(vec![
2212 Field::new("a", DataType::Int16, false),
2213 Field::new("b", DataType::Utf8, false),
2214 Field::new("c", DataType::Boolean, false),
2215 ]));
2216
2217 let batches = do_read(buf, 1024, true, true, schema);
2218 assert_eq!(batches.len(), 1);
2219
2220 let buf = r#"
2221 {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
2222 {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
2223 "#;
2224
2225 let schema = Arc::new(Schema::new(vec![
2226 Field::new("a", DataType::Int16, false),
2227 Field::new("b", DataType::Utf8, false),
2228 Field::new_struct(
2229 "c",
2230 vec![
2231 Field::new("a", DataType::Boolean, false),
2232 Field::new("b", DataType::Int16, false),
2233 ],
2234 false,
2235 ),
2236 ]));
2237
2238 let batches = do_read(buf, 1024, true, true, schema);
2239 assert_eq!(batches.len(), 1);
2240 }
2241
2242 #[test]
2243 fn test_strict_mode_missing_columns_in_schema() {
2244 let buf = r#"
2245 {"a": 1, "b": "2", "c": true}
2246 {"a": 2E0, "b": "4", "c": false}
2247 "#;
2248
2249 let schema = Arc::new(Schema::new(vec![
2250 Field::new("a", DataType::Int16, true),
2251 Field::new("c", DataType::Boolean, true),
2252 ]));
2253
2254 let err = ReaderBuilder::new(schema)
2255 .with_batch_size(1024)
2256 .with_strict_mode(true)
2257 .build(Cursor::new(buf.as_bytes()))
2258 .unwrap()
2259 .read()
2260 .unwrap_err();
2261
2262 assert_eq!(
2263 err.to_string(),
2264 "Json error: column 'b' missing from schema"
2265 );
2266
2267 let buf = r#"
2268 {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
2269 {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
2270 "#;
2271
2272 let schema = Arc::new(Schema::new(vec![
2273 Field::new("a", DataType::Int16, false),
2274 Field::new("b", DataType::Utf8, false),
2275 Field::new_struct("c", vec![Field::new("a", DataType::Boolean, false)], false),
2276 ]));
2277
2278 let err = ReaderBuilder::new(schema)
2279 .with_batch_size(1024)
2280 .with_strict_mode(true)
2281 .build(Cursor::new(buf.as_bytes()))
2282 .unwrap()
2283 .read()
2284 .unwrap_err();
2285
2286 assert_eq!(
2287 err.to_string(),
2288 "Json error: whilst decoding field 'c': column 'b' missing from schema"
2289 );
2290 }
2291
2292 fn read_file(path: &str, schema: Option<Schema>) -> Reader<BufReader<File>> {
2293 let file = File::open(path).unwrap();
2294 let mut reader = BufReader::new(file);
2295 let schema = schema.unwrap_or_else(|| {
2296 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2297 reader.rewind().unwrap();
2298 schema
2299 });
2300 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
2301 builder.build(reader).unwrap()
2302 }
2303
2304 #[test]
2305 fn test_json_basic() {
2306 let mut reader = read_file("test/data/basic.json", None);
2307 let batch = reader.next().unwrap().unwrap();
2308
2309 assert_eq!(8, batch.num_columns());
2310 assert_eq!(12, batch.num_rows());
2311
2312 let schema = reader.schema();
2313 let batch_schema = batch.schema();
2314 assert_eq!(schema, batch_schema);
2315
2316 let a = schema.column_with_name("a").unwrap();
2317 assert_eq!(0, a.0);
2318 assert_eq!(&DataType::Int64, a.1.data_type());
2319 let b = schema.column_with_name("b").unwrap();
2320 assert_eq!(1, b.0);
2321 assert_eq!(&DataType::Float64, b.1.data_type());
2322 let c = schema.column_with_name("c").unwrap();
2323 assert_eq!(2, c.0);
2324 assert_eq!(&DataType::Boolean, c.1.data_type());
2325 let d = schema.column_with_name("d").unwrap();
2326 assert_eq!(3, d.0);
2327 assert_eq!(&DataType::Utf8, d.1.data_type());
2328
2329 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2330 assert_eq!(1, aa.value(0));
2331 assert_eq!(-10, aa.value(1));
2332 let bb = batch.column(b.0).as_primitive::<Float64Type>();
2333 assert_eq!(2.0, bb.value(0));
2334 assert_eq!(-3.5, bb.value(1));
2335 let cc = batch.column(c.0).as_boolean();
2336 assert!(!cc.value(0));
2337 assert!(cc.value(10));
2338 let dd = batch.column(d.0).as_string::<i32>();
2339 assert_eq!("4", dd.value(0));
2340 assert_eq!("text", dd.value(8));
2341 }
2342
2343 #[test]
2344 fn test_json_empty_projection() {
2345 let mut reader = read_file("test/data/basic.json", Some(Schema::empty()));
2346 let batch = reader.next().unwrap().unwrap();
2347
2348 assert_eq!(0, batch.num_columns());
2349 assert_eq!(12, batch.num_rows());
2350 }
2351
2352 #[test]
2353 fn test_json_basic_with_nulls() {
2354 let mut reader = read_file("test/data/basic_nulls.json", None);
2355 let batch = reader.next().unwrap().unwrap();
2356
2357 assert_eq!(4, batch.num_columns());
2358 assert_eq!(12, batch.num_rows());
2359
2360 let schema = reader.schema();
2361 let batch_schema = batch.schema();
2362 assert_eq!(schema, batch_schema);
2363
2364 let a = schema.column_with_name("a").unwrap();
2365 assert_eq!(&DataType::Int64, a.1.data_type());
2366 let b = schema.column_with_name("b").unwrap();
2367 assert_eq!(&DataType::Float64, b.1.data_type());
2368 let c = schema.column_with_name("c").unwrap();
2369 assert_eq!(&DataType::Boolean, c.1.data_type());
2370 let d = schema.column_with_name("d").unwrap();
2371 assert_eq!(&DataType::Utf8, d.1.data_type());
2372
2373 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2374 assert!(aa.is_valid(0));
2375 assert!(!aa.is_valid(1));
2376 assert!(!aa.is_valid(11));
2377 let bb = batch.column(b.0).as_primitive::<Float64Type>();
2378 assert!(bb.is_valid(0));
2379 assert!(!bb.is_valid(2));
2380 assert!(!bb.is_valid(11));
2381 let cc = batch.column(c.0).as_boolean();
2382 assert!(cc.is_valid(0));
2383 assert!(!cc.is_valid(4));
2384 assert!(!cc.is_valid(11));
2385 let dd = batch.column(d.0).as_string::<i32>();
2386 assert!(!dd.is_valid(0));
2387 assert!(dd.is_valid(1));
2388 assert!(!dd.is_valid(4));
2389 assert!(!dd.is_valid(11));
2390 }
2391
2392 #[test]
2393 fn test_json_basic_schema() {
2394 let schema = Schema::new(vec![
2395 Field::new("a", DataType::Int64, true),
2396 Field::new("b", DataType::Float32, false),
2397 Field::new("c", DataType::Boolean, false),
2398 Field::new("d", DataType::Utf8, false),
2399 ]);
2400
2401 let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2402 let reader_schema = reader.schema();
2403 assert_eq!(reader_schema.as_ref(), &schema);
2404 let batch = reader.next().unwrap().unwrap();
2405
2406 assert_eq!(4, batch.num_columns());
2407 assert_eq!(12, batch.num_rows());
2408
2409 let schema = batch.schema();
2410
2411 let a = schema.column_with_name("a").unwrap();
2412 assert_eq!(&DataType::Int64, a.1.data_type());
2413 let b = schema.column_with_name("b").unwrap();
2414 assert_eq!(&DataType::Float32, b.1.data_type());
2415 let c = schema.column_with_name("c").unwrap();
2416 assert_eq!(&DataType::Boolean, c.1.data_type());
2417 let d = schema.column_with_name("d").unwrap();
2418 assert_eq!(&DataType::Utf8, d.1.data_type());
2419
2420 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2421 assert_eq!(1, aa.value(0));
2422 assert_eq!(100000000000000, aa.value(11));
2423 let bb = batch.column(b.0).as_primitive::<Float32Type>();
2424 assert_eq!(2.0, bb.value(0));
2425 assert_eq!(-3.5, bb.value(1));
2426 }
2427
2428 #[test]
2429 fn test_json_basic_schema_projection() {
2430 let schema = Schema::new(vec![
2431 Field::new("a", DataType::Int64, true),
2432 Field::new("c", DataType::Boolean, false),
2433 ]);
2434
2435 let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2436 let batch = reader.next().unwrap().unwrap();
2437
2438 assert_eq!(2, batch.num_columns());
2439 assert_eq!(2, batch.schema().fields().len());
2440 assert_eq!(12, batch.num_rows());
2441
2442 assert_eq!(batch.schema().as_ref(), &schema);
2443
2444 let a = schema.column_with_name("a").unwrap();
2445 assert_eq!(0, a.0);
2446 assert_eq!(&DataType::Int64, a.1.data_type());
2447 let c = schema.column_with_name("c").unwrap();
2448 assert_eq!(1, c.0);
2449 assert_eq!(&DataType::Boolean, c.1.data_type());
2450 }
2451
2452 #[test]
2453 fn test_json_arrays() {
2454 let mut reader = read_file("test/data/arrays.json", None);
2455 let batch = reader.next().unwrap().unwrap();
2456
2457 assert_eq!(4, batch.num_columns());
2458 assert_eq!(3, batch.num_rows());
2459
2460 let schema = batch.schema();
2461
2462 let a = schema.column_with_name("a").unwrap();
2463 assert_eq!(&DataType::Int64, a.1.data_type());
2464 let b = schema.column_with_name("b").unwrap();
2465 assert_eq!(
2466 &DataType::List(Arc::new(Field::new_list_field(DataType::Float64, true))),
2467 b.1.data_type()
2468 );
2469 let c = schema.column_with_name("c").unwrap();
2470 assert_eq!(
2471 &DataType::List(Arc::new(Field::new_list_field(DataType::Boolean, true))),
2472 c.1.data_type()
2473 );
2474 let d = schema.column_with_name("d").unwrap();
2475 assert_eq!(&DataType::Utf8, d.1.data_type());
2476
2477 let aa = batch.column(a.0).as_primitive::<Int64Type>();
2478 assert_eq!(1, aa.value(0));
2479 assert_eq!(-10, aa.value(1));
2480 assert_eq!(1627668684594000000, aa.value(2));
2481 let bb = batch.column(b.0).as_list::<i32>();
2482 let bb = bb.values().as_primitive::<Float64Type>();
2483 assert_eq!(9, bb.len());
2484 assert_eq!(2.0, bb.value(0));
2485 assert_eq!(-6.1, bb.value(5));
2486 assert!(!bb.is_valid(7));
2487
2488 let cc = batch
2489 .column(c.0)
2490 .as_any()
2491 .downcast_ref::<ListArray>()
2492 .unwrap();
2493 let cc = cc.values().as_boolean();
2494 assert_eq!(6, cc.len());
2495 assert!(!cc.value(0));
2496 assert!(!cc.value(4));
2497 assert!(!cc.is_valid(5));
2498 }
2499
2500 #[test]
2501 fn test_empty_json_arrays() {
2502 let json_content = r#"
2503 {"items": []}
2504 {"items": null}
2505 {}
2506 "#;
2507
2508 let schema = Arc::new(Schema::new(vec![Field::new(
2509 "items",
2510 DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2511 true,
2512 )]));
2513
2514 let batches = do_read(json_content, 1024, false, false, schema);
2515 assert_eq!(batches.len(), 1);
2516
2517 let col1 = batches[0].column(0).as_list::<i32>();
2518 assert_eq!(col1.null_count(), 2);
2519 assert!(col1.value(0).is_empty());
2520 assert_eq!(col1.value(0).data_type(), &DataType::Null);
2521 assert!(col1.is_null(1));
2522 assert!(col1.is_null(2));
2523 }
2524
2525 #[test]
2526 fn test_nested_empty_json_arrays() {
2527 let json_content = r#"
2528 {"items": [[],[]]}
2529 {"items": [[null, null],[null]]}
2530 "#;
2531
2532 let schema = Arc::new(Schema::new(vec![Field::new(
2533 "items",
2534 DataType::List(FieldRef::new(Field::new_list_field(
2535 DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2536 true,
2537 ))),
2538 true,
2539 )]));
2540
2541 let batches = do_read(json_content, 1024, false, false, schema);
2542 assert_eq!(batches.len(), 1);
2543
2544 let col1 = batches[0].column(0).as_list::<i32>();
2545 assert_eq!(col1.null_count(), 0);
2546 assert_eq!(col1.value(0).len(), 2);
2547 assert!(col1.value(0).as_list::<i32>().value(0).is_empty());
2548 assert!(col1.value(0).as_list::<i32>().value(1).is_empty());
2549
2550 assert_eq!(col1.value(1).len(), 2);
2551 assert_eq!(col1.value(1).as_list::<i32>().value(0).len(), 2);
2552 assert_eq!(col1.value(1).as_list::<i32>().value(1).len(), 1);
2553 }
2554
2555 #[test]
2556 fn test_nested_list_json_arrays() {
2557 let c_field = Field::new_struct("c", vec![Field::new("d", DataType::Utf8, true)], true);
2558 let a_struct_field = Field::new_struct(
2559 "a",
2560 vec![Field::new("b", DataType::Boolean, true), c_field.clone()],
2561 true,
2562 );
2563 let a_field = Field::new("a", DataType::List(Arc::new(a_struct_field.clone())), true);
2564 let schema = Arc::new(Schema::new(vec![a_field.clone()]));
2565 let builder = ReaderBuilder::new(schema).with_batch_size(64);
2566 let json_content = r#"
2567 {"a": [{"b": true, "c": {"d": "a_text"}}, {"b": false, "c": {"d": "b_text"}}]}
2568 {"a": [{"b": false, "c": null}]}
2569 {"a": [{"b": true, "c": {"d": "c_text"}}, {"b": null, "c": {"d": "d_text"}}, {"b": true, "c": {"d": null}}]}
2570 {"a": null}
2571 {"a": []}
2572 {"a": [null]}
2573 "#;
2574 let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2575
2576 let d = StringArray::from(vec![
2578 Some("a_text"),
2579 Some("b_text"),
2580 None,
2581 Some("c_text"),
2582 Some("d_text"),
2583 None,
2584 None,
2585 ]);
2586 let c = StructArray::new(
2587 vec![Field::new("d", DataType::Utf8, true)].into(),
2588 vec![Arc::new(d.clone()) as ArrayRef],
2589 Some(NullBuffer::from(vec![
2590 true, true, false, true, true, true, false,
2591 ])),
2592 );
2593 let b = BooleanArray::from(vec![
2594 Some(true),
2595 Some(false),
2596 Some(false),
2597 Some(true),
2598 None,
2599 Some(true),
2600 None,
2601 ]);
2602 let a = StructArray::new(
2603 vec![Field::new("b", DataType::Boolean, true), c_field.clone()].into(),
2604 vec![
2605 Arc::new(b.clone()) as ArrayRef,
2606 Arc::new(c.clone()) as ArrayRef,
2607 ],
2608 Some(NullBuffer::from(vec![
2609 true, true, true, true, true, true, false,
2610 ])),
2611 );
2612 let a_list = ListArray::new(
2613 Arc::new(a_struct_field.clone()),
2614 OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 6, 6, 6, 7])),
2615 Arc::new(a),
2616 Some(NullBuffer::from(vec![true, true, true, false, true, true])),
2617 );
2618
2619 let batch = reader.next().unwrap().unwrap();
2621 let read = batch.column(0);
2622 assert_eq!(read.len(), 6);
2623 let read: &ListArray = read.as_list::<i32>();
2625 let expected = &a_list;
2626 assert_eq!(read.value_offsets(), &[0, 2, 3, 6, 6, 6, 7]);
2627 assert_eq!(read.nulls(), expected.nulls());
2629 let struct_array = read.values().as_struct();
2631 let expected_struct_array = expected.values().as_struct();
2632
2633 assert_eq!(7, struct_array.len());
2634 assert_eq!(1, struct_array.null_count());
2635 assert_eq!(7, expected_struct_array.len());
2636 assert_eq!(1, expected_struct_array.null_count());
2637 assert_eq!(struct_array.nulls(), expected_struct_array.nulls());
2639 let read_b = struct_array.column(0);
2641 assert_eq!(read_b.as_ref(), &b);
2642 let read_c = struct_array.column(1);
2643 assert_eq!(read_c.as_struct(), &c);
2644 let read_c = read_c.as_struct();
2645 let read_d = read_c.column(0);
2646 assert_eq!(read_d.as_ref(), &d);
2647
2648 assert_eq!(read, expected);
2649 }
2650
2651 fn assert_read_list_view<O: OffsetSizeTrait>() {
2652 let field = Arc::new(Field::new("item", DataType::Int32, true));
2653 let data_type = GenericListViewArray::<O>::DATA_TYPE_CONSTRUCTOR(field.clone());
2654 let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2655
2656 let buf = r#"
2657 {"lv": [1, 2, 3]}
2658 {"lv": [4, null]}
2659 {"lv": null}
2660 {"lv": [6]}
2661 {"lv": []}
2662 "#;
2663
2664 let batches = do_read(buf, 1024, false, false, schema);
2665 assert_eq!(batches.len(), 1);
2666 let batch = &batches[0];
2667 let col = batch.column(0);
2668 let list_view = col
2669 .as_any()
2670 .downcast_ref::<GenericListViewArray<O>>()
2671 .unwrap();
2672
2673 assert_eq!(list_view.len(), 5);
2674
2675 let expected_offsets: Vec<O> = vec![0, 3, 5, 5, 6]
2677 .into_iter()
2678 .map(|v| O::usize_as(v))
2679 .collect();
2680 let expected_sizes: Vec<O> = vec![3, 2, 0, 1, 0]
2681 .into_iter()
2682 .map(|v| O::usize_as(v))
2683 .collect();
2684 assert_eq!(list_view.value_offsets(), &expected_offsets);
2685 assert_eq!(list_view.value_sizes(), &expected_sizes);
2686
2687 assert!(list_view.is_valid(0));
2689 let vals = list_view.value(0);
2690 let ints = vals.as_primitive::<Int32Type>();
2691 assert_eq!(ints.values(), &[1, 2, 3]);
2692
2693 assert!(list_view.is_valid(1));
2695 let vals = list_view.value(1);
2696 let ints = vals.as_primitive::<Int32Type>();
2697 assert_eq!(ints.len(), 2);
2698 assert_eq!(ints.value(0), 4);
2699 assert!(ints.is_null(1));
2700
2701 assert!(list_view.is_null(2));
2703
2704 assert!(list_view.is_valid(3));
2706 let vals = list_view.value(3);
2707 let ints = vals.as_primitive::<Int32Type>();
2708 assert_eq!(ints.values(), &[6]);
2709
2710 assert!(list_view.is_valid(4));
2712 let vals = list_view.value(4);
2713 assert_eq!(vals.len(), 0);
2714 }
2715
2716 #[test]
2717 fn test_read_list_view() {
2718 assert_read_list_view::<i32>();
2719 assert_read_list_view::<i64>();
2720 }
2721
2722 #[test]
2723 fn test_read_list_view_rejects_null_non_nullable_child() {
2724 let field = Arc::new(Field::new("item", DataType::Int32, false));
2725 for (data_type, array_type) in [
2726 (DataType::ListView(field.clone()), "ListViewArray"),
2727 (DataType::LargeListView(field.clone()), "LargeListViewArray"),
2728 ] {
2729 let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2730 let buf = r#"
2731 {"lv": [1, 2, 3]}
2732 {"lv": [4, null]}
2733 "#;
2734
2735 let error = ReaderBuilder::new(schema)
2736 .build(Cursor::new(buf.as_bytes()))
2737 .unwrap()
2738 .collect::<Result<Vec<_>, _>>()
2739 .unwrap_err();
2740
2741 assert_eq!(
2742 error.to_string(),
2743 format!(
2744 "Invalid argument error: Non-nullable field of {array_type} \"item\" cannot contain nulls"
2745 )
2746 );
2747 }
2748 }
2749
2750 #[test]
2751 fn test_fixed_size_list() {
2752 let buf = r#"
2753 {"a": [1, 2, 3]}
2754 {"a": [4, 5, 6]}
2755 {"a": [7, 8, 9]}
2756 "#;
2757
2758 let field = Field::new_list_field(DataType::Int32, true);
2759 let schema = Arc::new(Schema::new(vec![Field::new(
2760 "a",
2761 DataType::FixedSizeList(Arc::new(field), 3),
2762 false,
2763 )]));
2764
2765 let batches = do_read(buf, 1024, false, false, schema);
2766 assert_eq!(batches.len(), 1);
2767
2768 let col = batches[0].column(0).as_fixed_size_list();
2769 assert_eq!(col.len(), 3);
2770 assert_eq!(col.value_length(), 3);
2771
2772 let values = col.values().as_primitive::<Int32Type>();
2773 assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8, 9]);
2774 }
2775
2776 #[test]
2777 fn test_fixed_size_list_nullable() {
2778 let buf = r#"
2779 {"a": [1, 2]}
2780 {"a": null}
2781 {"a": [3, null]}
2782 "#;
2783
2784 let field = Field::new_list_field(DataType::Int32, true);
2785 let schema = Arc::new(Schema::new(vec![Field::new(
2786 "a",
2787 DataType::FixedSizeList(Arc::new(field), 2),
2788 true,
2789 )]));
2790
2791 let batches = do_read(buf, 1024, false, false, schema);
2792 assert_eq!(batches.len(), 1);
2793
2794 let col = batches[0].column(0).as_fixed_size_list();
2795 assert_eq!(col.len(), 3);
2796 assert!(col.is_valid(0));
2797 assert!(col.is_null(1));
2798 assert!(col.is_valid(2));
2799
2800 let values = col.values().as_primitive::<Int32Type>();
2801 assert_eq!(values.value(0), 1);
2802 assert_eq!(values.value(1), 2);
2803 assert_eq!(values.value(4), 3);
2804 assert!(values.is_null(5));
2805 }
2806
2807 #[test]
2808 fn test_fixed_size_list_zero_size_non_nullable() {
2809 let buf = r#"
2810 {"a": []}
2811 {"a": []}
2812 {"a": []}
2813 "#;
2814
2815 let field = Field::new_list_field(DataType::Int32, true);
2816 let schema = Arc::new(Schema::new(vec![Field::new(
2817 "a",
2818 DataType::FixedSizeList(Arc::new(field), 0),
2819 false,
2820 )]));
2821
2822 let batches = do_read(buf, 1024, false, false, schema);
2823 assert_eq!(batches.len(), 1);
2824
2825 let col = batches[0].column(0).as_fixed_size_list();
2826 assert_eq!(col.len(), 3);
2827 assert_eq!(col.value_length(), 0);
2828
2829 let values = col.values().as_primitive::<Int32Type>();
2830 assert!(values.values().is_empty());
2831 }
2832
2833 #[test]
2834 fn test_fixed_size_list_wrong_size() {
2835 let buf = r#"{"a": [1, 2, 3]}"#;
2836
2837 let field = Field::new_list_field(DataType::Int32, true);
2838 let schema = Arc::new(Schema::new(vec![Field::new(
2839 "a",
2840 DataType::FixedSizeList(Arc::new(field), 2),
2841 false,
2842 )]));
2843
2844 let err = ReaderBuilder::new(schema)
2845 .build(Cursor::new(buf.as_bytes()))
2846 .unwrap()
2847 .next()
2848 .unwrap()
2849 .unwrap_err();
2850
2851 assert!(err.to_string().contains("expected 2 but got 3"), "{}", err);
2852 }
2853
2854 #[test]
2855 fn test_fixed_size_list_nested() {
2856 let buf = r#"
2857 {"a": [[1, 2], [3, 4]]}
2858 {"a": [[5, 6], [7, 8]]}
2859 "#;
2860
2861 let inner_field = Field::new_list_field(DataType::Int32, true);
2862 let inner_type = DataType::FixedSizeList(Arc::new(inner_field), 2);
2863 let outer_field = Arc::new(Field::new_list_field(inner_type.clone(), true));
2864 let schema = Arc::new(Schema::new(vec![Field::new(
2865 "a",
2866 DataType::FixedSizeList(outer_field, 2),
2867 false,
2868 )]));
2869
2870 let batches = do_read(buf, 1024, false, false, schema);
2871 assert_eq!(batches.len(), 1);
2872
2873 let col = batches[0].column(0).as_fixed_size_list();
2874 assert_eq!(col.len(), 2);
2875 assert_eq!(col.value_length(), 2);
2876
2877 let inner = col.values().as_fixed_size_list();
2878 assert_eq!(inner.len(), 4);
2879 assert_eq!(inner.value_length(), 2);
2880
2881 let values = inner.values().as_primitive::<Int32Type>();
2882 assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8]);
2883 }
2884
2885 #[test]
2886 fn test_fixed_size_list_ignore_type_conflicts() {
2887 let field = Field::new("item", DataType::Int32, true);
2888 let schema = Arc::new(Schema::new(vec![Field::new(
2889 "a",
2890 DataType::FixedSizeList(Arc::new(field), 2),
2891 true,
2892 )]));
2893
2894 let json = vec![
2895 json!({"a": [1, 2]}),
2896 json!({"a": "not a list"}),
2897 json!({"a": 42}),
2898 json!({"a": [6, 7]}),
2899 ];
2900
2901 let mut decoder = ReaderBuilder::new(schema)
2902 .with_ignore_type_conflicts(true)
2903 .build_decoder()
2904 .unwrap();
2905 decoder.serialize(&json).unwrap();
2906 let batch = decoder.flush().unwrap().unwrap();
2907
2908 let col = batch.column(0).as_fixed_size_list();
2909 assert_eq!(col.len(), 4);
2910 assert!(col.is_valid(0));
2911 assert!(col.is_null(1)); assert!(col.is_null(2)); assert!(col.is_valid(3));
2914
2915 let values = col.values().as_primitive::<Int32Type>();
2916 assert_eq!(values.value(0), 1);
2917 assert_eq!(values.value(1), 2);
2918 assert_eq!(values.value(6), 6);
2919 assert_eq!(values.value(7), 7);
2920 }
2921
2922 #[test]
2923 fn test_skip_empty_lines() {
2924 let schema = Schema::new(vec![Field::new("a", DataType::Int64, true)]);
2925 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
2926 let json_content = "
2927 {\"a\": 1}
2928 {\"a\": 2}
2929 {\"a\": 3}";
2930 let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2931 let batch = reader.next().unwrap().unwrap();
2932
2933 assert_eq!(1, batch.num_columns());
2934 assert_eq!(3, batch.num_rows());
2935
2936 let schema = reader.schema();
2937 let c = schema.column_with_name("a").unwrap();
2938 assert_eq!(&DataType::Int64, c.1.data_type());
2939 }
2940
2941 #[test]
2942 fn test_with_multiple_batches() {
2943 let file = File::open("test/data/basic_nulls.json").unwrap();
2944 let mut reader = BufReader::new(file);
2945 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2946 reader.rewind().unwrap();
2947
2948 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2949 let mut reader = builder.build(reader).unwrap();
2950
2951 let mut num_records = Vec::new();
2952 while let Some(rb) = reader.next().transpose().unwrap() {
2953 num_records.push(rb.num_rows());
2954 }
2955
2956 assert_eq!(vec![5, 5, 2], num_records);
2957 }
2958
2959 #[test]
2960 fn test_timestamp_from_json_seconds() {
2961 let schema = Schema::new(vec![Field::new(
2962 "a",
2963 DataType::Timestamp(TimeUnit::Second, None),
2964 true,
2965 )]);
2966
2967 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2968 let batch = reader.next().unwrap().unwrap();
2969
2970 assert_eq!(1, batch.num_columns());
2971 assert_eq!(12, batch.num_rows());
2972
2973 let schema = reader.schema();
2974 let batch_schema = batch.schema();
2975 assert_eq!(schema, batch_schema);
2976
2977 let a = schema.column_with_name("a").unwrap();
2978 assert_eq!(
2979 &DataType::Timestamp(TimeUnit::Second, None),
2980 a.1.data_type()
2981 );
2982
2983 let aa = batch.column(a.0).as_primitive::<TimestampSecondType>();
2984 assert!(aa.is_valid(0));
2985 assert!(!aa.is_valid(1));
2986 assert!(!aa.is_valid(2));
2987 assert_eq!(1, aa.value(0));
2988 assert_eq!(1, aa.value(3));
2989 assert_eq!(5, aa.value(7));
2990 }
2991
2992 #[test]
2993 fn test_timestamp_from_json_milliseconds() {
2994 let schema = Schema::new(vec![Field::new(
2995 "a",
2996 DataType::Timestamp(TimeUnit::Millisecond, None),
2997 true,
2998 )]);
2999
3000 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
3001 let batch = reader.next().unwrap().unwrap();
3002
3003 assert_eq!(1, batch.num_columns());
3004 assert_eq!(12, batch.num_rows());
3005
3006 let schema = reader.schema();
3007 let batch_schema = batch.schema();
3008 assert_eq!(schema, batch_schema);
3009
3010 let a = schema.column_with_name("a").unwrap();
3011 assert_eq!(
3012 &DataType::Timestamp(TimeUnit::Millisecond, None),
3013 a.1.data_type()
3014 );
3015
3016 let aa = batch.column(a.0).as_primitive::<TimestampMillisecondType>();
3017 assert!(aa.is_valid(0));
3018 assert!(!aa.is_valid(1));
3019 assert!(!aa.is_valid(2));
3020 assert_eq!(1, aa.value(0));
3021 assert_eq!(1, aa.value(3));
3022 assert_eq!(5, aa.value(7));
3023 }
3024
3025 #[test]
3026 fn test_date_from_json_milliseconds() {
3027 let schema = Schema::new(vec![Field::new("a", DataType::Date64, true)]);
3028
3029 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
3030 let batch = reader.next().unwrap().unwrap();
3031
3032 assert_eq!(1, batch.num_columns());
3033 assert_eq!(12, batch.num_rows());
3034
3035 let schema = reader.schema();
3036 let batch_schema = batch.schema();
3037 assert_eq!(schema, batch_schema);
3038
3039 let a = schema.column_with_name("a").unwrap();
3040 assert_eq!(&DataType::Date64, a.1.data_type());
3041
3042 let aa = batch.column(a.0).as_primitive::<Date64Type>();
3043 assert!(aa.is_valid(0));
3044 assert!(!aa.is_valid(1));
3045 assert!(!aa.is_valid(2));
3046 assert_eq!(1, aa.value(0));
3047 assert_eq!(1, aa.value(3));
3048 assert_eq!(5, aa.value(7));
3049 }
3050
3051 #[test]
3052 fn test_time_from_json_nanoseconds() {
3053 let schema = Schema::new(vec![Field::new(
3054 "a",
3055 DataType::Time64(TimeUnit::Nanosecond),
3056 true,
3057 )]);
3058
3059 let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
3060 let batch = reader.next().unwrap().unwrap();
3061
3062 assert_eq!(1, batch.num_columns());
3063 assert_eq!(12, batch.num_rows());
3064
3065 let schema = reader.schema();
3066 let batch_schema = batch.schema();
3067 assert_eq!(schema, batch_schema);
3068
3069 let a = schema.column_with_name("a").unwrap();
3070 assert_eq!(&DataType::Time64(TimeUnit::Nanosecond), a.1.data_type());
3071
3072 let aa = batch.column(a.0).as_primitive::<Time64NanosecondType>();
3073 assert!(aa.is_valid(0));
3074 assert!(!aa.is_valid(1));
3075 assert!(!aa.is_valid(2));
3076 assert_eq!(1, aa.value(0));
3077 assert_eq!(1, aa.value(3));
3078 assert_eq!(5, aa.value(7));
3079 }
3080
3081 #[test]
3082 fn test_json_iterator() {
3083 let file = File::open("test/data/basic.json").unwrap();
3084 let mut reader = BufReader::new(file);
3085 let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
3086 reader.rewind().unwrap();
3087
3088 let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
3089 let reader = builder.build(reader).unwrap();
3090 let schema = reader.schema();
3091 let (col_a_index, _) = schema.column_with_name("a").unwrap();
3092
3093 let mut sum_num_rows = 0;
3094 let mut num_batches = 0;
3095 let mut sum_a = 0;
3096 for batch in reader {
3097 let batch = batch.unwrap();
3098 assert_eq!(8, batch.num_columns());
3099 sum_num_rows += batch.num_rows();
3100 num_batches += 1;
3101 let batch_schema = batch.schema();
3102 assert_eq!(schema, batch_schema);
3103 let a_array = batch.column(col_a_index).as_primitive::<Int64Type>();
3104 sum_a += (0..a_array.len()).map(|i| a_array.value(i)).sum::<i64>();
3105 }
3106 assert_eq!(12, sum_num_rows);
3107 assert_eq!(3, num_batches);
3108 assert_eq!(100000000000011, sum_a);
3109 }
3110
3111 #[test]
3112 fn test_decoder_error() {
3113 let schema = Arc::new(Schema::new(vec![Field::new_struct(
3114 "a",
3115 vec![Field::new("child", DataType::Int32, false)],
3116 true,
3117 )]));
3118
3119 let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
3120 let _ = decoder.decode(br#"{"a": { "child":"#).unwrap();
3121 assert!(decoder.tape_decoder.has_partial_row());
3122 assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
3123 let _ = decoder.flush().unwrap_err();
3124 assert!(decoder.tape_decoder.has_partial_row());
3125 assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
3126
3127 let parse_err = |s: &str| {
3128 ReaderBuilder::new(schema.clone())
3129 .build(Cursor::new(s.as_bytes()))
3130 .unwrap()
3131 .next()
3132 .unwrap()
3133 .unwrap_err()
3134 .to_string()
3135 };
3136
3137 let err = parse_err(r#"{"a": 123}"#);
3138 assert_eq!(
3139 err,
3140 "Json error: whilst decoding field 'a': expected { got 123"
3141 );
3142
3143 let err = parse_err(r#"{"a": ["bar"]}"#);
3144 assert_eq!(
3145 err,
3146 r#"Json error: whilst decoding field 'a': expected { got ["bar"]"#
3147 );
3148
3149 let err = parse_err(r#"{"a": []}"#);
3150 assert_eq!(
3151 err,
3152 "Json error: whilst decoding field 'a': expected { got []"
3153 );
3154
3155 let err = parse_err(r#"{"a": [{"child": 234}]}"#);
3156 assert_eq!(
3157 err,
3158 r#"Json error: whilst decoding field 'a': expected { got [{"child": 234}]"#
3159 );
3160
3161 let err = parse_err(r#"{"a": [{"child": {"foo": [{"foo": ["bar"]}]}}]}"#);
3162 assert_eq!(
3163 err,
3164 r#"Json error: whilst decoding field 'a': expected { got [{"child": {"foo": [{"foo": ["bar"]}]}}]"#
3165 );
3166
3167 let err = parse_err(r#"{"a": true}"#);
3168 assert_eq!(
3169 err,
3170 "Json error: whilst decoding field 'a': expected { got true"
3171 );
3172
3173 let err = parse_err(r#"{"a": false}"#);
3174 assert_eq!(
3175 err,
3176 "Json error: whilst decoding field 'a': expected { got false"
3177 );
3178
3179 let err = parse_err(r#"{"a": "foo"}"#);
3180 assert_eq!(
3181 err,
3182 "Json error: whilst decoding field 'a': expected { got \"foo\""
3183 );
3184
3185 let err = parse_err(r#"{"a": {"child": false}}"#);
3186 assert_eq!(
3187 err,
3188 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got false"
3189 );
3190
3191 let err = parse_err(r#"{"a": {"child": []}}"#);
3192 assert_eq!(
3193 err,
3194 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got []"
3195 );
3196
3197 let err = parse_err(r#"{"a": {"child": [123]}}"#);
3198 assert_eq!(
3199 err,
3200 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123]"
3201 );
3202
3203 let err = parse_err(r#"{"a": {"child": [123, 3465346]}}"#);
3204 assert_eq!(
3205 err,
3206 "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123, 3465346]"
3207 );
3208 }
3209
3210 #[test]
3211 fn test_serialize_timestamp() {
3212 let json = vec![
3213 json!({"timestamp": 1681319393}),
3214 json!({"timestamp": "1970-01-01T00:00:00+02:00"}),
3215 ];
3216 let schema = Schema::new(vec![Field::new(
3217 "timestamp",
3218 DataType::Timestamp(TimeUnit::Second, None),
3219 true,
3220 )]);
3221 let mut decoder = ReaderBuilder::new(Arc::new(schema))
3222 .build_decoder()
3223 .unwrap();
3224 decoder.serialize(&json).unwrap();
3225 let batch = decoder.flush().unwrap().unwrap();
3226 assert_eq!(batch.num_rows(), 2);
3227 assert_eq!(batch.num_columns(), 1);
3228 let values = batch.column(0).as_primitive::<TimestampSecondType>();
3229 assert_eq!(values.values(), &[1681319393, -7200]);
3230 }
3231
3232 #[test]
3233 fn test_serialize_decimal() {
3234 let json = vec![
3235 json!({"decimal": 1.234}),
3236 json!({"decimal": "1.234"}),
3237 json!({"decimal": 1234}),
3238 json!({"decimal": "1234"}),
3239 ];
3240 let schema = Schema::new(vec![Field::new(
3241 "decimal",
3242 DataType::Decimal128(10, 3),
3243 true,
3244 )]);
3245 let mut decoder = ReaderBuilder::new(Arc::new(schema))
3246 .build_decoder()
3247 .unwrap();
3248 decoder.serialize(&json).unwrap();
3249 let batch = decoder.flush().unwrap().unwrap();
3250 assert_eq!(batch.num_rows(), 4);
3251 assert_eq!(batch.num_columns(), 1);
3252 let values = batch.column(0).as_primitive::<Decimal128Type>();
3253 assert_eq!(values.values(), &[1234, 1234, 1234000, 1234000]);
3254 }
3255
3256 #[test]
3257 fn test_serde_field() {
3258 let field = Field::new("int", DataType::Int32, true);
3259 let mut decoder = ReaderBuilder::new_with_field(field)
3260 .build_decoder()
3261 .unwrap();
3262 decoder.serialize(&[1_i32, 2, 3, 4]).unwrap();
3263 let b = decoder.flush().unwrap().unwrap();
3264 let values = b.column(0).as_primitive::<Int32Type>().values();
3265 assert_eq!(values, &[1, 2, 3, 4]);
3266 }
3267
3268 #[test]
3269 fn test_serde_large_numbers() {
3270 let field = Field::new("int", DataType::Int64, true);
3271 let mut decoder = ReaderBuilder::new_with_field(field)
3272 .build_decoder()
3273 .unwrap();
3274
3275 decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
3276 let b = decoder.flush().unwrap().unwrap();
3277 let values = b.column(0).as_primitive::<Int64Type>().values();
3278 assert_eq!(values, &[1699148028689, 2, 3, 4]);
3279
3280 let field = Field::new(
3281 "int",
3282 DataType::Timestamp(TimeUnit::Microsecond, None),
3283 true,
3284 );
3285 let mut decoder = ReaderBuilder::new_with_field(field)
3286 .build_decoder()
3287 .unwrap();
3288
3289 decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
3290 let b = decoder.flush().unwrap().unwrap();
3291 let values = b
3292 .column(0)
3293 .as_primitive::<TimestampMicrosecondType>()
3294 .values();
3295 assert_eq!(values, &[1699148028689, 2, 3, 4]);
3296 }
3297
3298 #[test]
3299 fn test_coercing_primitive_into_string_decoder() {
3300 let buf = &format!(
3301 r#"[{{"a": 1, "b": "A", "c": "T"}}, {{"a": 2, "b": "BB", "c": "F"}}, {{"a": {}, "b": 123, "c": false}}, {{"a": {}, "b": 789, "c": true}}]"#,
3302 (i32::MAX as i64 + 10),
3303 i64::MAX - 10
3304 );
3305 let schema = Schema::new(vec![
3306 Field::new("a", DataType::Float64, true),
3307 Field::new("b", DataType::Utf8, true),
3308 Field::new("c", DataType::Utf8, true),
3309 ]);
3310 let json_array: Vec<serde_json::Value> = serde_json::from_str(buf).unwrap();
3311 let schema_ref = Arc::new(schema);
3312
3313 let reader = ReaderBuilder::new(schema_ref.clone()).with_coerce_primitive(true);
3315 let mut decoder = reader.build_decoder().unwrap();
3316 decoder.serialize(json_array.as_slice()).unwrap();
3317 let batch = decoder.flush().unwrap().unwrap();
3318 assert_eq!(
3319 batch,
3320 RecordBatch::try_new(
3321 schema_ref,
3322 vec![
3323 Arc::new(Float64Array::from(vec![
3324 1.0,
3325 2.0,
3326 (i32::MAX as i64 + 10) as f64,
3327 (i64::MAX - 10) as f64
3328 ])),
3329 Arc::new(StringArray::from(vec!["A", "BB", "123", "789"])),
3330 Arc::new(StringArray::from(vec!["T", "F", "false", "true"])),
3331 ]
3332 )
3333 .unwrap()
3334 );
3335 }
3336
3337 #[test]
3338 fn test_serialize_f32_into_string() {
3339 let field = Field::new("f", DataType::Utf8, true);
3341 let mut decoder = ReaderBuilder::new_with_field(field)
3342 .with_coerce_primitive(true)
3343 .build_decoder()
3344 .unwrap();
3345 decoder.serialize(&[1.5_f32, -2.25_f32]).unwrap();
3346 let batch = decoder.flush().unwrap().unwrap();
3347 let values = batch.column(0).as_string::<i32>();
3348 assert_eq!(values.value(0), "1.5");
3349 assert_eq!(values.value(1), "-2.25");
3350 }
3351
3352 fn _parse_structs(
3357 row: &str,
3358 struct_mode: StructMode,
3359 fields: Fields,
3360 as_struct: bool,
3361 ) -> Result<RecordBatch, ArrowError> {
3362 let builder = if as_struct {
3363 ReaderBuilder::new_with_field(Field::new("r", DataType::Struct(fields), true))
3364 } else {
3365 ReaderBuilder::new(Arc::new(Schema::new(fields)))
3366 };
3367 builder
3368 .with_struct_mode(struct_mode)
3369 .build(Cursor::new(row.as_bytes()))
3370 .unwrap()
3371 .next()
3372 .unwrap()
3373 }
3374
3375 #[test]
3376 fn test_struct_decoding_list_length() {
3377 use arrow_array::array;
3378
3379 let row = "[1, 2]";
3380
3381 let mut fields = vec![Field::new("a", DataType::Int32, true)];
3382 let too_few_fields = Fields::from(fields.clone());
3383 fields.push(Field::new("b", DataType::Int32, true));
3384 let correct_fields = Fields::from(fields.clone());
3385 fields.push(Field::new("c", DataType::Int32, true));
3386 let too_many_fields = Fields::from(fields.clone());
3387
3388 let parse = |fields: Fields, as_struct: bool| {
3389 _parse_structs(row, StructMode::ListOnly, fields, as_struct)
3390 };
3391
3392 let expected_row = StructArray::new(
3393 correct_fields.clone(),
3394 vec![
3395 Arc::new(array::Int32Array::from(vec![1])),
3396 Arc::new(array::Int32Array::from(vec![2])),
3397 ],
3398 None,
3399 );
3400 let row_field = Field::new("r", DataType::Struct(correct_fields.clone()), true);
3401
3402 assert_eq!(
3403 parse(too_few_fields.clone(), true).unwrap_err().to_string(),
3404 "Json error: found extra columns for 1 fields".to_string()
3405 );
3406 assert_eq!(
3407 parse(too_few_fields, false).unwrap_err().to_string(),
3408 "Json error: found extra columns for 1 fields".to_string()
3409 );
3410 assert_eq!(
3411 parse(correct_fields.clone(), true).unwrap(),
3412 RecordBatch::try_new(
3413 Arc::new(Schema::new(vec![row_field])),
3414 vec![Arc::new(expected_row.clone())]
3415 )
3416 .unwrap()
3417 );
3418 assert_eq!(
3419 parse(correct_fields, false).unwrap(),
3420 RecordBatch::from(expected_row)
3421 );
3422 assert_eq!(
3423 parse(too_many_fields.clone(), true)
3424 .unwrap_err()
3425 .to_string(),
3426 "Json error: found 2 columns for 3 fields".to_string()
3427 );
3428 assert_eq!(
3429 parse(too_many_fields, false).unwrap_err().to_string(),
3430 "Json error: found 2 columns for 3 fields".to_string()
3431 );
3432 }
3433
3434 #[test]
3435 fn test_struct_decoding() {
3436 use arrow_array::builder;
3437
3438 let nested_object_json = r#"{"a": {"b": [1, 2], "c": {"d": 3}}}"#;
3439 let nested_list_json = r#"[[[1, 2], {"d": 3}]]"#;
3440 let nested_mixed_json = r#"{"a": [[1, 2], {"d": 3}]}"#;
3441
3442 let struct_fields = Fields::from(vec![
3443 Field::new("b", DataType::new_list(DataType::Int32, true), true),
3444 Field::new_map(
3445 "c",
3446 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3447 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3448 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Int32, true),
3449 false,
3450 false,
3451 ),
3452 ]);
3453
3454 let list_array =
3455 ListArray::from_iter_primitive::<Int32Type, _, _>(vec![Some(vec![Some(1), Some(2)])]);
3456
3457 let map_array = {
3458 let mut map_builder = builder::MapBuilder::new(
3459 None,
3460 builder::StringBuilder::new(),
3461 builder::Int32Builder::new(),
3462 );
3463 map_builder.keys().append_value("d");
3464 map_builder.values().append_value(3);
3465 map_builder.append(true).unwrap();
3466 map_builder.finish()
3467 };
3468
3469 let struct_array = StructArray::new(
3470 struct_fields.clone(),
3471 vec![Arc::new(list_array), Arc::new(map_array)],
3472 None,
3473 );
3474
3475 let fields = Fields::from(vec![Field::new("a", DataType::Struct(struct_fields), true)]);
3476 let schema = Arc::new(Schema::new(fields.clone()));
3477 let expected = RecordBatch::try_new(schema.clone(), vec![Arc::new(struct_array)]).unwrap();
3478
3479 let parse = |row: &str, struct_mode: StructMode| {
3480 _parse_structs(row, struct_mode, fields.clone(), false)
3481 };
3482
3483 assert_eq!(
3484 parse(nested_object_json, StructMode::ObjectOnly).unwrap(),
3485 expected
3486 );
3487 assert_eq!(
3488 parse(nested_list_json, StructMode::ObjectOnly)
3489 .unwrap_err()
3490 .to_string(),
3491 "Json error: expected { got [[[1, 2], {\"d\": 3}]]".to_owned()
3492 );
3493 assert_eq!(
3494 parse(nested_mixed_json, StructMode::ObjectOnly)
3495 .unwrap_err()
3496 .to_string(),
3497 "Json error: whilst decoding field 'a': expected { got [[1, 2], {\"d\": 3}]".to_owned()
3498 );
3499
3500 assert_eq!(
3501 parse(nested_list_json, StructMode::ListOnly).unwrap(),
3502 expected
3503 );
3504 assert_eq!(
3505 parse(nested_object_json, StructMode::ListOnly)
3506 .unwrap_err()
3507 .to_string(),
3508 "Json error: expected [ got {\"a\": {\"b\": [1, 2]\"c\": {\"d\": 3}}}".to_owned()
3509 );
3510 assert_eq!(
3511 parse(nested_mixed_json, StructMode::ListOnly)
3512 .unwrap_err()
3513 .to_string(),
3514 "Json error: expected [ got {\"a\": [[1, 2], {\"d\": 3}]}".to_owned()
3515 );
3516 }
3517
3518 #[test]
3524 fn test_struct_decoding_empty_list() {
3525 let int_field = Field::new("a", DataType::Int32, true);
3526 let struct_field = Field::new(
3527 "r",
3528 DataType::Struct(Fields::from(vec![int_field.clone()])),
3529 true,
3530 );
3531
3532 let parse = |row: &str, as_struct: bool, field: Field| {
3533 _parse_structs(
3534 row,
3535 StructMode::ListOnly,
3536 Fields::from(vec![field]),
3537 as_struct,
3538 )
3539 };
3540
3541 assert_eq!(
3543 parse("[]", true, struct_field.clone())
3544 .unwrap_err()
3545 .to_string(),
3546 "Json error: found 0 columns for 1 fields".to_owned()
3547 );
3548 assert_eq!(
3549 parse("[]", false, int_field.clone())
3550 .unwrap_err()
3551 .to_string(),
3552 "Json error: found 0 columns for 1 fields".to_owned()
3553 );
3554 assert_eq!(
3555 parse("[]", false, struct_field.clone())
3556 .unwrap_err()
3557 .to_string(),
3558 "Json error: found 0 columns for 1 fields".to_owned()
3559 );
3560 assert_eq!(
3561 parse("[[]]", false, struct_field.clone())
3562 .unwrap_err()
3563 .to_string(),
3564 "Json error: whilst decoding field 'r': found 0 columns for 1 fields".to_owned()
3565 );
3566 }
3567
3568 #[test]
3569 fn test_decode_list_struct_with_wrong_types() {
3570 let int_field = Field::new("a", DataType::Int32, true);
3571 let struct_field = Field::new(
3572 "r",
3573 DataType::Struct(Fields::from(vec![int_field.clone()])),
3574 true,
3575 );
3576
3577 let parse = |row: &str, as_struct: bool, field: Field| {
3578 _parse_structs(
3579 row,
3580 StructMode::ListOnly,
3581 Fields::from(vec![field]),
3582 as_struct,
3583 )
3584 };
3585
3586 assert_eq!(
3588 parse(r#"[["a"]]"#, false, struct_field.clone())
3589 .unwrap_err()
3590 .to_string(),
3591 "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3592 );
3593 assert_eq!(
3594 parse(r#"[["a"]]"#, true, struct_field.clone())
3595 .unwrap_err()
3596 .to_string(),
3597 "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3598 );
3599 assert_eq!(
3600 parse(r#"["a"]"#, true, int_field.clone())
3601 .unwrap_err()
3602 .to_string(),
3603 "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3604 );
3605 assert_eq!(
3606 parse(r#"["a"]"#, false, int_field.clone())
3607 .unwrap_err()
3608 .to_string(),
3609 "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3610 );
3611 }
3612
3613 #[test]
3614 fn test_type_conflict_nulls() {
3615 let schema = Schema::new(vec![
3616 Field::new("null", DataType::Null, true),
3617 Field::new("bool", DataType::Boolean, true),
3618 Field::new("primitive", DataType::Int32, true),
3619 Field::new("numeric", DataType::Decimal128(10, 3), true),
3620 Field::new("string", DataType::Utf8, true),
3621 Field::new("string_view", DataType::Utf8View, true),
3622 Field::new(
3623 "timestamp",
3624 DataType::Timestamp(TimeUnit::Second, None),
3625 true,
3626 ),
3627 Field::new(
3628 "array",
3629 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3630 true,
3631 ),
3632 Field::new(
3633 "map",
3634 DataType::Map(
3635 Arc::new(Field::new(
3636 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3637 DataType::Struct(Fields::from(vec![
3638 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3639 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3640 ])),
3641 false, )),
3643 false, ),
3645 true, ),
3647 Field::new(
3648 "struct",
3649 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3650 true,
3651 ),
3652 ]);
3653
3654 let json_values = vec![
3656 json!(null),
3657 json!(true),
3658 json!(42),
3659 json!(1.234),
3660 json!("hi"),
3661 json!("ho"),
3662 json!("1970-01-01T00:00:00+02:00"),
3663 json!([1, "ho", 3]),
3664 json!({"k": "value"}),
3665 json!({"a": 1}),
3666 ];
3667
3668 let json: Vec<_> = (0..json_values.len())
3670 .map(|i| {
3671 let pairs = json_values[i..]
3672 .iter()
3673 .chain(json_values[..i].iter())
3674 .zip(&schema.fields)
3675 .map(|(v, f)| (f.name().clone(), v.clone()))
3676 .collect();
3677 serde_json::Value::Object(pairs)
3678 })
3679 .collect();
3680 let mut decoder = ReaderBuilder::new(Arc::new(schema))
3681 .with_ignore_type_conflicts(true)
3682 .with_coerce_primitive(true)
3683 .build_decoder()
3684 .unwrap();
3685 decoder.serialize(&json).unwrap();
3686 let batch = decoder.flush().unwrap().unwrap();
3687 assert_eq!(batch.num_rows(), 10);
3688 assert_eq!(batch.num_columns(), 10);
3689
3690 let _ = batch
3692 .column(0)
3693 .as_any()
3694 .downcast_ref::<NullArray>()
3695 .unwrap();
3696
3697 assert!(
3698 batch
3699 .column(1)
3700 .as_any()
3701 .downcast_ref::<BooleanArray>()
3702 .unwrap()
3703 .iter()
3704 .eq([
3705 Some(true),
3706 None,
3707 None,
3708 None,
3709 None,
3710 None,
3711 None,
3712 None,
3713 None,
3714 None
3715 ])
3716 );
3717
3718 assert!(batch.column(2).as_primitive::<Int32Type>().iter().eq([
3719 Some(42),
3720 Some(1),
3721 None,
3722 None,
3723 None,
3724 None,
3725 None,
3726 None,
3727 None,
3728 None
3729 ]));
3730
3731 assert!(batch.column(3).as_primitive::<Decimal128Type>().iter().eq([
3732 Some(1234),
3733 None,
3734 None,
3735 None,
3736 None,
3737 None,
3738 None,
3739 None,
3740 None,
3741 Some(42000)
3742 ]));
3743
3744 assert!(
3745 batch
3746 .column(4)
3747 .as_any()
3748 .downcast_ref::<StringArray>()
3749 .unwrap()
3750 .iter()
3751 .eq([
3752 Some("hi"),
3753 Some("ho"),
3754 Some("1970-01-01T00:00:00+02:00"),
3755 None,
3756 None,
3757 None,
3758 None,
3759 Some("true"),
3760 Some("42"),
3761 Some("1.234"),
3762 ])
3763 );
3764
3765 assert!(
3766 batch
3767 .column(5)
3768 .as_any()
3769 .downcast_ref::<StringViewArray>()
3770 .unwrap()
3771 .iter()
3772 .eq([
3773 Some("ho"),
3774 Some("1970-01-01T00:00:00+02:00"),
3775 None,
3776 None,
3777 None,
3778 None,
3779 Some("true"),
3780 Some("42"),
3781 Some("1.234"),
3782 Some("hi"),
3783 ])
3784 );
3785
3786 assert!(
3787 batch
3788 .column(6)
3789 .as_primitive::<TimestampSecondType>()
3790 .iter()
3791 .eq([
3792 Some(-7200),
3793 None,
3794 None,
3795 None,
3796 None,
3797 None,
3798 Some(42),
3799 None,
3800 None,
3801 None,
3802 ])
3803 );
3804
3805 let arrays = batch
3806 .column(7)
3807 .as_any()
3808 .downcast_ref::<ListArray>()
3809 .unwrap();
3810 assert_eq!(
3811 arrays.nulls(),
3812 Some(&NullBuffer::from(
3813 &[
3814 true, false, false, false, false, false, false, false, false, false
3815 ][..]
3816 ))
3817 );
3818 assert_eq!(arrays.offsets()[1], 3);
3819 let array_values = arrays
3820 .values()
3821 .as_any()
3822 .downcast_ref::<Int32Array>()
3823 .unwrap();
3824 assert!(array_values.iter().eq([Some(1), None, Some(3)]));
3825
3826 let maps = batch.column(8).as_any().downcast_ref::<MapArray>().unwrap();
3827 assert_eq!(
3828 maps.nulls(),
3829 Some(&NullBuffer::from(
3830 &[
3832 true, true, false, false, false, false, false, false, false, false
3833 ][..]
3834 ))
3835 );
3836 let map_keys = maps.keys().as_any().downcast_ref::<StringArray>().unwrap();
3837 assert!(map_keys.iter().eq([Some("k"), Some("a")]));
3838 let map_values = maps
3839 .values()
3840 .as_any()
3841 .downcast_ref::<StringArray>()
3842 .unwrap();
3843 assert!(map_values.iter().eq([Some("value"), Some("1")]));
3844
3845 let structs = batch
3846 .column(9)
3847 .as_any()
3848 .downcast_ref::<StructArray>()
3849 .unwrap();
3850 assert_eq!(
3851 structs.nulls(),
3852 Some(&NullBuffer::from(
3853 &[
3855 true, false, false, false, false, false, false, false, false, true
3856 ][..]
3857 ))
3858 );
3859 let struct_fields = structs
3860 .column(0)
3861 .as_any()
3862 .downcast_ref::<Int32Array>()
3863 .unwrap();
3864 assert!(struct_fields.slice(0, 2).iter().eq([Some(1), None]));
3865 }
3866
3867 #[test]
3868 fn test_type_conflict_non_nullable() {
3869 let fields = [
3870 Field::new("bool", DataType::Boolean, false),
3871 Field::new("primitive", DataType::Int32, false),
3872 Field::new("numeric", DataType::Decimal128(10, 3), false),
3873 Field::new("string", DataType::Utf8, false),
3874 Field::new("string_view", DataType::Utf8View, false),
3875 Field::new(
3876 "timestamp",
3877 DataType::Timestamp(TimeUnit::Second, None),
3878 false,
3879 ),
3880 Field::new(
3881 "array",
3882 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3883 false,
3884 ),
3885 Field::new(
3886 "fixed_size_list",
3887 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3888 false,
3889 ),
3890 Field::new(
3891 "map",
3892 DataType::Map(
3893 Arc::new(Field::new(
3894 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3895 DataType::Struct(Fields::from(vec![
3896 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3897 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3898 ])),
3899 false, )),
3901 false, ),
3903 false, ),
3905 Field::new(
3906 "struct",
3907 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3908 false,
3909 ),
3910 ];
3911
3912 let json_values = vec![json!(true), json!({"a": 1})];
3914
3915 for field in fields {
3916 let mut decoder = ReaderBuilder::new_with_field(field)
3917 .with_ignore_type_conflicts(true)
3918 .build_decoder()
3919 .unwrap();
3920 decoder.serialize(&json_values).unwrap();
3921 decoder
3922 .flush()
3923 .expect_err("type conflict on non-nullable type");
3924 }
3925 }
3926
3927 #[test]
3928 fn test_ignore_type_conflicts_disabled() {
3929 let fields = [
3930 Field::new("null", DataType::Null, true),
3931 Field::new("bool", DataType::Boolean, true),
3932 Field::new("primitive", DataType::Int32, true),
3933 Field::new("numeric", DataType::Decimal128(10, 3), true),
3934 Field::new("string", DataType::Utf8, true),
3935 Field::new("string_view", DataType::Utf8View, true),
3936 Field::new(
3937 "timestamp",
3938 DataType::Timestamp(TimeUnit::Second, None),
3939 true,
3940 ),
3941 Field::new(
3942 "array",
3943 DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3944 true,
3945 ),
3946 Field::new(
3947 "fixed_size_list",
3948 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3949 true,
3950 ),
3951 Field::new(
3952 "map",
3953 DataType::Map(
3954 Arc::new(Field::new(
3955 Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3956 DataType::Struct(Fields::from(vec![
3957 Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3958 Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3959 ])),
3960 false, )),
3962 false, ),
3964 true, ),
3966 Field::new(
3967 "struct",
3968 DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3969 true,
3970 ),
3971 ];
3972
3973 let json_values = vec![json!(true), json!({"a": 1})];
3975
3976 for field in fields {
3977 let mut decoder = ReaderBuilder::new_with_field(field)
3978 .build_decoder()
3979 .unwrap();
3980 decoder.serialize(&json_values).unwrap();
3981 decoder
3982 .flush()
3983 .expect_err("type conflict on non-nullable type");
3984 }
3985 }
3986
3987 #[test]
3988 fn test_read_run_end_encoded() {
3989 let buf = r#"
3990 {"a": "x"}
3991 {"a": "x"}
3992 {"a": "y"}
3993 {"a": "y"}
3994 {"a": "y"}
3995 "#;
3996
3997 let ree_type = DataType::RunEndEncoded(
3998 Arc::new(Field::new("run_ends", DataType::Int32, false)),
3999 Arc::new(Field::new("values", DataType::Utf8, true)),
4000 );
4001 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
4002 let batches = do_read(buf, 1024, false, false, schema);
4003 assert_eq!(batches.len(), 1);
4004
4005 let col = batches[0].column(0);
4006 let run_array = col.as_run::<arrow_array::types::Int32Type>();
4007
4008 assert_eq!(run_array.len(), 5);
4010 assert_eq!(run_array.run_ends().values(), &[2, 5]);
4011
4012 let values = run_array.values().as_string::<i32>();
4013 assert_eq!(values.len(), 2);
4014 assert_eq!(values.value(0), "x");
4015 assert_eq!(values.value(1), "y");
4016 }
4017
4018 #[test]
4019 fn test_read_run_end_encoded_consecutive_nulls() {
4020 let buf = r#"
4021 {"a": "x"}
4022 {}
4023 {}
4024 {}
4025 {"a": "y"}
4026 "#;
4027
4028 let ree_type = DataType::RunEndEncoded(
4029 Arc::new(Field::new("run_ends", DataType::Int32, false)),
4030 Arc::new(Field::new("values", DataType::Utf8, true)),
4031 );
4032 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
4033 let batches = do_read(buf, 1024, false, false, schema);
4034 assert_eq!(batches.len(), 1);
4035
4036 let col = batches[0].column(0);
4037 let run_array = col.as_run::<arrow_array::types::Int32Type>();
4038
4039 assert_eq!(run_array.len(), 5);
4041 assert_eq!(run_array.run_ends().values(), &[1, 4, 5]);
4042
4043 let values = run_array.values().as_string::<i32>();
4044 assert_eq!(values.len(), 3);
4045 assert_eq!(values.value(0), "x");
4046 assert!(values.is_null(1));
4047 assert_eq!(values.value(2), "y");
4048 }
4049
4050 #[test]
4051 fn test_read_run_end_encoded_nullability() {
4052 for field_nullable in [false, true] {
4053 for values_nullable in [false, true] {
4054 let ree_type = DataType::RunEndEncoded(
4055 Arc::new(Field::new("run_ends", DataType::Int32, false)),
4056 Arc::new(Field::new("values", DataType::Utf8, values_nullable)),
4057 );
4058 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, field_nullable)]));
4059
4060 for buf in [
4061 r#"{"a": "x"}
4062 {"a": null}
4063 {"a": "y"}"#,
4064 r#"{"a": "x"}
4065 {}
4066 {"a": "y"}"#,
4067 ] {
4068 let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
4069 let result = decoder.decode(buf.as_bytes()).and_then(|_| decoder.flush());
4070
4071 if field_nullable && values_nullable {
4072 result.expect("REE field and values are both nullable");
4073 } else {
4074 let err = result.expect_err(
4075 "REE nulls require both the field and values to be nullable",
4076 );
4077 assert!(
4078 err.to_string().contains(
4079 "Encountered nulls in non-nullable values of RunEndEncoded"
4080 ),
4081 "unexpected error: {err}"
4082 );
4083 }
4084 }
4085 }
4086 }
4087 }
4088
4089 #[test]
4090 fn test_read_run_end_encoded_all_unique() {
4091 let buf = r#"
4092 {"a": 1}
4093 {"a": 2}
4094 {"a": 3}
4095 "#;
4096
4097 let ree_type = DataType::RunEndEncoded(
4098 Arc::new(Field::new("run_ends", DataType::Int32, false)),
4099 Arc::new(Field::new("values", DataType::Int32, true)),
4100 );
4101 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
4102 let batches = do_read(buf, 1024, false, false, schema);
4103 assert_eq!(batches.len(), 1);
4104
4105 let col = batches[0].column(0);
4106 let run_array = col.as_run::<arrow_array::types::Int32Type>();
4107
4108 assert_eq!(run_array.len(), 3);
4110 assert_eq!(run_array.run_ends().values(), &[1, 2, 3]);
4111 }
4112
4113 #[test]
4114 fn test_read_run_end_encoded_int16_run_ends() {
4115 let buf = r#"
4116 {"a": "x"}
4117 {"a": "x"}
4118 {"a": "y"}
4119 "#;
4120
4121 let ree_type = DataType::RunEndEncoded(
4122 Arc::new(Field::new("run_ends", DataType::Int16, false)),
4123 Arc::new(Field::new("values", DataType::Utf8, true)),
4124 );
4125 let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
4126 let batches = do_read(buf, 1024, false, false, schema);
4127 assert_eq!(batches.len(), 1);
4128
4129 let col = batches[0].column(0);
4130 let run_array = col.as_run::<arrow_array::types::Int16Type>();
4131
4132 assert_eq!(run_array.len(), 3);
4133 assert_eq!(run_array.run_ends().values(), &[2i16, 3]);
4134 }
4135
4136 #[test]
4137 fn test_read_nested_run_end_encoded() {
4138 let buf = r#"
4139 {"a": "x"}
4140 {"a": "x"}
4141 {"a": "y"}
4142 "#;
4143
4144 let inner_type = DataType::RunEndEncoded(
4147 Arc::new(Field::new("run_ends", DataType::Int64, false)),
4148 Arc::new(Field::new("values", DataType::Utf8, true)),
4149 );
4150 let outer_type = DataType::RunEndEncoded(
4151 Arc::new(Field::new("run_ends", DataType::Int64, false)),
4152 Arc::new(Field::new("values", inner_type, true)),
4153 );
4154 let schema = Arc::new(Schema::new(vec![Field::new("a", outer_type, true)]));
4155 let batches = do_read(buf, 1024, false, false, schema);
4156 assert_eq!(batches.len(), 1);
4157
4158 let col = batches[0].column(0);
4159 let outer = col.as_run::<arrow_array::types::Int64Type>();
4160 assert_eq!(outer.len(), 3);
4162 assert_eq!(outer.run_ends().values(), &[2, 3]);
4163
4164 let nested = outer.values().as_run::<arrow_array::types::Int64Type>();
4165 assert_eq!(nested.len(), 2);
4167 assert_eq!(nested.run_ends().values(), &[1, 2]);
4168
4169 let nested_values = nested.values().as_string::<i32>();
4170 assert_eq!(nested_values.len(), 2);
4171 assert_eq!(nested_values.value(0), "x");
4172 assert_eq!(nested_values.value(1), "y");
4173 }
4174
4175 #[test]
4176 fn test_flatten_top_level_arrays() {
4177 let buf = r#"
4178 [
4179 {"a": 1},
4180 {"a": 2}
4181 ]
4182 {"a": 3}
4183 [{"a": 4}, {"a": 5}, {"a": 6}]
4184 "#;
4185
4186 let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, true)]));
4187 let batches = do_read_config(buf, ReaderBuilder::new(schema).with_flatten(true));
4188 assert_eq!(batches.len(), 1);
4189
4190 let col = batches[0].column(0);
4191 let col = col.as_any().downcast_ref::<Int32Array>().unwrap();
4192 assert_eq!(col.len(), 6);
4193 for i in 0..6 {
4194 assert_eq!(col.value(i), (i as i32) + 1);
4195 }
4196 }
4197
4198 #[test]
4200 fn test_decoder_factory_declining_everything_is_transparent() {
4201 #[derive(Debug, Default)]
4202 struct DeclineAll {
4203 seen: Mutex<Vec<FieldRef>>,
4204 }
4205
4206 impl DecoderFactory for DeclineAll {
4207 fn make_default_decoder(
4208 &self,
4209 _ctx: &DecoderContext,
4210 field: &FieldRef,
4211 _is_nullable: bool,
4212 ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> {
4213 self.seen.lock().unwrap().push(field.clone());
4214 Ok(None)
4215 }
4216 }
4217
4218 let schema = Arc::new(Schema::new(vec![
4219 Field::new("a", DataType::Int32, true),
4220 Field::new("b", DataType::Utf8, true),
4221 Field::new_list("c", Field::new("item", DataType::Float64, true), true),
4222 ]));
4223 let buf = r#"{"a": 1, "b": "x", "c": [1.5, 2.5]}
4224 {"a": null, "c": []}
4225 "#;
4226
4227 let read = |factory: Option<Arc<dyn DecoderFactory>>| {
4228 let mut builder = ReaderBuilder::new(schema.clone());
4229 if let Some(factory) = factory {
4230 builder = builder.with_decoder_factory(factory);
4231 }
4232 builder
4233 .build(Cursor::new(buf.as_bytes()))
4234 .unwrap()
4235 .next()
4236 .unwrap()
4237 .unwrap()
4238 };
4239
4240 let factory = Arc::new(DeclineAll::default());
4241 assert_eq!(read(None), read(Some(factory.clone())));
4242
4243 let seen = factory.seen.lock().unwrap();
4246 assert!(
4247 matches!(seen[0].data_type(), DataType::Struct(_)),
4248 "{seen:?}"
4249 );
4250 assert_eq!(seen[0].name(), "", "root field should be nameless");
4251
4252 let by_name: Vec<(&str, &DataType)> = seen
4253 .iter()
4254 .map(|f| (f.name().as_str(), f.data_type()))
4255 .collect();
4256 for expected in [
4257 ("a", &DataType::Int32),
4258 ("b", &DataType::Utf8),
4259 ("item", &DataType::Float64),
4260 ] {
4261 assert!(
4262 by_name.contains(&expected),
4263 "{expected:?} missing from {by_name:?}"
4264 );
4265 }
4266 }
4267
4268 #[test]
4271 fn test_decoder_factory_output_is_validated() {
4272 struct Bad {
4274 wrong_type: bool,
4275 }
4276
4277 impl ArrayDecoder for Bad {
4278 fn decode(&mut self, _tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
4279 Ok(match self.wrong_type {
4280 true => Arc::new(StringViewArray::from(vec!["x"; pos.len()])) as ArrayRef,
4281 false => Arc::new(StringArray::from(Vec::<&str>::new())),
4282 })
4283 }
4284 }
4285
4286 #[derive(Debug)]
4287 struct BadFactory {
4288 wrong_type: bool,
4289 }
4290
4291 impl DecoderFactory for BadFactory {
4292 fn make_default_decoder(
4293 &self,
4294 _ctx: &DecoderContext,
4295 field: &FieldRef,
4296 _is_nullable: bool,
4297 ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> {
4298 match field.data_type() {
4299 DataType::Utf8 => Ok(Some(Box::new(Bad {
4300 wrong_type: self.wrong_type,
4301 }))),
4302 _ => Ok(None),
4303 }
4304 }
4305 }
4306
4307 let read = |wrong_type: bool| {
4308 let schema = Arc::new(Schema::new(vec![Field::new("s", DataType::Utf8, true)]));
4309 ReaderBuilder::new(schema)
4310 .with_decoder_factory(Arc::new(BadFactory { wrong_type }))
4311 .build(Cursor::new(br#"{"s": "hello"}"#.as_slice()))
4312 .unwrap()
4313 .next()
4314 .unwrap()
4315 .unwrap_err()
4316 .to_string()
4317 };
4318
4319 let err = read(true);
4320 assert!(
4321 err.contains("returned Utf8View for a field of type Utf8"),
4322 "{err}"
4323 );
4324
4325 let err = read(false);
4326 assert!(err.contains("returned 0 values for 1 rows"), "{err}");
4327 }
4328}