Skip to main content

arrow_json/reader/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! JSON reader
19//!
20//! This JSON reader allows JSON records to be read into the Arrow memory
21//! model. Records are loaded in batches and are then converted from the record-oriented
22//! representation to the columnar arrow data model.
23//!
24//! The reader ignores whitespace between JSON values, including `\n` and `\r`, allowing
25//! parsing of sequences of one or more arbitrarily formatted JSON values, including
26//! but not limited to newline-delimited JSON.
27//!
28//! # Basic Usage
29//!
30//! [`Reader`] can be used directly with synchronous data sources, such as [`std::fs::File`]
31//!
32//! ```
33//! # use arrow_schema::*;
34//! # use std::fs::File;
35//! # use std::io::BufReader;
36//! # use std::sync::Arc;
37//!
38//! let schema = Arc::new(Schema::new(vec![
39//!     Field::new("a", DataType::Float64, false),
40//!     Field::new("b", DataType::Float64, false),
41//!     Field::new("c", DataType::Boolean, true),
42//! ]));
43//!
44//! let file = File::open("test/data/basic.json").unwrap();
45//!
46//! let mut json = arrow_json::ReaderBuilder::new(schema).build(BufReader::new(file)).unwrap();
47//! let batch = json.next().unwrap().unwrap();
48//! ```
49//!
50//! # Async Usage
51//!
52//! The lower-level [`Decoder`] can be integrated with various forms of async data streams,
53//! and is designed to be agnostic to the various different kinds of async IO primitives found
54//! within the Rust ecosystem.
55//!
56//! For example, see below for how it can be used with an arbitrary `Stream` of `Bytes`
57//!
58//! ```
59//! # use std::task::{Poll, ready};
60//! # use bytes::{Buf, Bytes};
61//! # use arrow_schema::ArrowError;
62//! # use futures::stream::{Stream, StreamExt};
63//! # use arrow_array::RecordBatch;
64//! # use arrow_json::reader::Decoder;
65//! #
66//! fn decode_stream<S: Stream<Item = Bytes> + Unpin>(
67//!     mut decoder: Decoder,
68//!     mut input: S,
69//! ) -> impl Stream<Item = Result<RecordBatch, ArrowError>> {
70//!     let mut buffered = Bytes::new();
71//!     futures::stream::poll_fn(move |cx| {
72//!         loop {
73//!             if buffered.is_empty() {
74//!                 buffered = match ready!(input.poll_next_unpin(cx)) {
75//!                     Some(b) => b,
76//!                     None => break,
77//!                 };
78//!             }
79//!             let decoded = match decoder.decode(buffered.as_ref()) {
80//!                 Ok(decoded) => decoded,
81//!                 Err(e) => return Poll::Ready(Some(Err(e))),
82//!             };
83//!             let read = buffered.len();
84//!             buffered.advance(decoded);
85//!             if decoded != read {
86//!                 break
87//!             }
88//!         }
89//!
90//!         Poll::Ready(decoder.flush().transpose())
91//!     })
92//! }
93//!
94//! ```
95//!
96//! In a similar vein, it can also be used with tokio-based IO primitives
97//!
98//! ```
99//! # use std::sync::Arc;
100//! # use arrow_schema::{DataType, Field, Schema};
101//! # use std::pin::Pin;
102//! # use std::task::{Poll, ready};
103//! # use futures::{Stream, TryStreamExt};
104//! # use tokio::io::AsyncBufRead;
105//! # use arrow_array::RecordBatch;
106//! # use arrow_json::reader::Decoder;
107//! # use arrow_schema::ArrowError;
108//! fn decode_stream<R: AsyncBufRead + Unpin>(
109//!     mut decoder: Decoder,
110//!     mut reader: R,
111//! ) -> impl Stream<Item = Result<RecordBatch, ArrowError>> {
112//!     futures::stream::poll_fn(move |cx| {
113//!         loop {
114//!             let b = match ready!(Pin::new(&mut reader).poll_fill_buf(cx)) {
115//!                 Ok(b) if b.is_empty() => break,
116//!                 Ok(b) => b,
117//!                 Err(e) => return Poll::Ready(Some(Err(e.into()))),
118//!             };
119//!             let read = b.len();
120//!             let decoded = match decoder.decode(b) {
121//!                 Ok(decoded) => decoded,
122//!                 Err(e) => return Poll::Ready(Some(Err(e))),
123//!             };
124//!             Pin::new(&mut reader).consume(decoded);
125//!             if decoded != read {
126//!                 break;
127//!             }
128//!         }
129//!
130//!         Poll::Ready(decoder.flush().transpose())
131//!     })
132//! }
133//! ```
134//!
135
136use 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::*;
167// `mod tape` stays private so `TapeDecoder` does not become public API
168pub 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/// A builder for [`Reader`] and [`Decoder`]
189#[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    /// Create a new [`ReaderBuilder`] with the provided [`SchemaRef`]
204    ///
205    /// This could be obtained using [`infer_json_schema`] if not known
206    ///
207    /// Any columns not present in `schema` will be ignored, unless `strict_mode` is set to true.
208    /// In this case, an error is returned when a column is missing from `schema`.
209    ///
210    /// [`infer_json_schema`]: crate::reader::infer_json_schema
211    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    /// Create a new [`ReaderBuilder`] that will parse JSON values of `field.data_type()`
226    ///
227    /// Unlike [`ReaderBuilder::new`] this does not require the root of the JSON data
228    /// to be an object, i.e. `{..}`, allowing for parsing of any valid JSON value(s)
229    ///
230    /// ```
231    /// # use std::sync::Arc;
232    /// # use arrow_array::cast::AsArray;
233    /// # use arrow_array::types::Int32Type;
234    /// # use arrow_json::ReaderBuilder;
235    /// # use arrow_schema::{DataType, Field};
236    /// // Root of JSON schema is a numeric type
237    /// let data = "1\n2\n3\n";
238    /// let field = Arc::new(Field::new("int", DataType::Int32, true));
239    /// let mut reader = ReaderBuilder::new_with_field(field.clone()).build(data.as_bytes()).unwrap();
240    /// let b = reader.next().unwrap().unwrap();
241    /// let values = b.column(0).as_primitive::<Int32Type>().values();
242    /// assert_eq!(values, &[1, 2, 3]);
243    ///
244    /// // Root of JSON schema is a list type
245    /// let data = "[1, 2, 3, 4, 5, 6, 7]\n[1, 2, 3]";
246    /// let field = Field::new_list("int", field.clone(), true);
247    /// let mut reader = ReaderBuilder::new_with_field(field).build(data.as_bytes()).unwrap();
248    /// let b = reader.next().unwrap().unwrap();
249    /// let list = b.column(0).as_list::<i32>();
250    ///
251    /// assert_eq!(list.offsets().as_ref(), &[0, 7, 10]);
252    /// let list_values = list.values().as_primitive::<Int32Type>();
253    /// assert_eq!(list_values.values(), &[1, 2, 3, 4, 5, 6, 7, 1, 2, 3]);
254    /// ```
255    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    /// Sets the batch size in rows to read
270    pub fn with_batch_size(self, batch_size: usize) -> Self {
271        Self { batch_size, ..self }
272    }
273
274    /// Sets if the decoder should coerce primitive values (bool and number) into string
275    /// when the Schema's column is Utf8 or LargeUtf8.
276    pub fn with_coerce_primitive(self, coerce_primitive: bool) -> Self {
277        Self {
278            coerce_primitive,
279            ..self
280        }
281    }
282
283    /// Sets if the decoder should return an error if it encounters a column not
284    /// present in `schema`. If `struct_mode` is `ListOnly` the value of
285    /// `strict_mode` is effectively `true`. It is required for all fields of
286    /// the struct to be in the list: without field names, there is no way to
287    /// determine which field is missing.
288    pub fn with_strict_mode(self, strict_mode: bool) -> Self {
289        Self {
290            strict_mode,
291            ..self
292        }
293    }
294
295    /// Set the [`StructMode`] for the reader, which determines whether structs
296    /// can be decoded from JSON as objects or lists. For more details refer to
297    /// the enum documentation. Default is to use `ObjectOnly`.
298    pub fn with_struct_mode(self, struct_mode: StructMode) -> Self {
299        Self {
300            struct_mode,
301            ..self
302        }
303    }
304
305    /// Sets whether the decoder should produce NULL instead of returning an error if it encounters
306    /// value that can not be parsed into the specified column type.
307    ///
308    /// For example, if the type is declared to be a nullable array of `DataType::Int32` but the
309    /// reader encounters a string value `"foo"` and the value `ignore_type_conflicts` is:
310    ///
311    /// * `false` (the default): The reader will return an error.
312    ///
313    /// * `true`: The reader will fill in NULL value for that array element.
314    ///
315    /// NOTE: An inferred NULL due to a type conflict will still produce parsing errors for
316    /// non-nullable fields, the same as any other NULL or missing value.
317    pub fn with_ignore_type_conflicts(self, ignore_type_conflicts: bool) -> Self {
318        Self {
319            ignore_type_conflicts,
320            ..self
321        }
322    }
323    /// Sets whether to flatten top-level arrays.
324    ///
325    /// * When `true`, each element of a top-level array will be treated as its own row.
326    /// * When `false` (the default), the entire top-level array will be treated as one row.
327    ///
328    /// For example, consider this input file:
329    /// ```text
330    /// [{ "a": 1 }, { "a": 2 }, { "b": 3 }]
331    /// [{ "a": 4 }, { "a": 5 }, { "b": 6 }]
332    /// ```
333    ///
334    /// By default, this would be parsed as two rows, each an array containing three elements.
335    /// With this option set to `true`, however, this would be parsed as six rows.
336    ///
337    /// Note that even with this option set to `true`, top-level objects are still permitted
338    /// and will be parsed as individual rows in the exact same manner as when this option is
339    /// set to `false`. It is even possible to mix top-level arrays with top-level objects.
340    pub fn with_flatten(self, flatten_top_level_arrays: bool) -> Self {
341        Self {
342            flatten_top_level_arrays,
343            ..self
344        }
345    }
346
347    /// Set a hook for customizing decoding behavior. See [`DecoderFactory`].
348    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    /// Create a [`Reader`] with the provided [`BufRead`]
356    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    /// Create a [`Decoder`]
364    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            // The root has no field of its own; synthesize a nameless one
371            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
400/// Reads JSON data with a known schema directly into arrow [`RecordBatch`]
401///
402/// Lines consisting solely of ASCII whitespace are ignored
403pub 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    /// Reads the next [`RecordBatch`] returning `Ok(None)` if EOF
418    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
450/// A low-level interface for reading JSON data from a byte stream
451///
452/// See [`Reader`] for a higher-level interface for interface with [`BufRead`]
453///
454/// The push-based interface facilitates integration with sources that yield arbitrarily
455/// delimited bytes ranges, such as [`BufRead`], or a chunked byte stream received from
456/// object storage
457///
458/// ```
459/// # use std::io::BufRead;
460/// # use arrow_array::RecordBatch;
461/// # use arrow_json::reader::{Decoder, ReaderBuilder};
462/// # use arrow_schema::{ArrowError, SchemaRef};
463/// #
464/// fn read_from_json<R: BufRead>(
465///     mut reader: R,
466///     schema: SchemaRef,
467/// ) -> Result<impl Iterator<Item = Result<RecordBatch, ArrowError>>, ArrowError> {
468///     let mut decoder = ReaderBuilder::new(schema).build_decoder()?;
469///     let mut next = move || {
470///         loop {
471///             // Decoder is agnostic that buf doesn't contain whole records
472///             let buf = reader.fill_buf()?;
473///             if buf.is_empty() {
474///                 break; // Input exhausted
475///             }
476///             let read = buf.len();
477///             let decoded = decoder.decode(buf)?;
478///
479///             // Consume the number of bytes read
480///             reader.consume(decoded);
481///             if decoded != read {
482///                 break; // Read batch size
483///             }
484///         }
485///         decoder.flush()
486///     };
487///     Ok(std::iter::from_fn(move || next().transpose()))
488/// }
489/// ```
490pub 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    /// Read JSON objects from `buf`, returning the number of bytes read
509    ///
510    /// This method returns once `batch_size` objects have been parsed since the
511    /// last call to [`Self::flush`], or `buf` is exhausted. Any remaining bytes
512    /// should be included in the next call to [`Self::decode`]
513    ///
514    /// There is no requirement that `buf` contains a whole number of records, facilitating
515    /// integration with arbitrary byte streams, such as those yielded by [`BufRead`]
516    pub fn decode(&mut self, buf: &[u8]) -> Result<usize, ArrowError> {
517        self.tape_decoder.decode(buf)
518    }
519
520    /// Serialize `rows` to this [`Decoder`]
521    ///
522    /// This provides a simple way to convert [serde]-compatible datastructures into arrow
523    /// [`RecordBatch`].
524    ///
525    /// Custom conversion logic as described in [arrow_array::builder] will likely outperform this,
526    /// especially where the schema is known at compile-time, however, this provides a mechanism
527    /// to get something up and running quickly
528    ///
529    /// It can be used with [`serde_json::Value`]
530    ///
531    /// ```
532    /// # use std::sync::Arc;
533    /// # use serde_json::{Value, json};
534    /// # use arrow_array::cast::AsArray;
535    /// # use arrow_array::types::Float32Type;
536    /// # use arrow_json::ReaderBuilder;
537    /// # use arrow_schema::{DataType, Field, Schema};
538    /// let json = vec![json!({"float": 2.3}), json!({"float": 5.7})];
539    ///
540    /// let schema = Schema::new(vec![Field::new("float", DataType::Float32, true)]);
541    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
542    ///
543    /// decoder.serialize(&json).unwrap();
544    /// let batch = decoder.flush().unwrap().unwrap();
545    /// assert_eq!(batch.num_rows(), 2);
546    /// assert_eq!(batch.num_columns(), 1);
547    /// let values = batch.column(0).as_primitive::<Float32Type>().values();
548    /// assert_eq!(values, &[2.3, 5.7])
549    /// ```
550    ///
551    /// Or with arbitrary [`Serialize`] types
552    ///
553    /// ```
554    /// # use std::sync::Arc;
555    /// # use arrow_json::ReaderBuilder;
556    /// # use arrow_schema::{DataType, Field, Schema};
557    /// # use serde::Serialize;
558    /// # use arrow_array::cast::AsArray;
559    /// # use arrow_array::types::{Float32Type, Int32Type};
560    /// #
561    /// #[derive(Serialize)]
562    /// struct MyStruct {
563    ///     int32: i32,
564    ///     float: f32,
565    /// }
566    ///
567    /// let schema = Schema::new(vec![
568    ///     Field::new("int32", DataType::Int32, false),
569    ///     Field::new("float", DataType::Float32, false),
570    /// ]);
571    ///
572    /// let rows = vec![
573    ///     MyStruct{ int32: 0, float: 3. },
574    ///     MyStruct{ int32: 4, float: 67.53 },
575    /// ];
576    ///
577    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
578    /// decoder.serialize(&rows).unwrap();
579    ///
580    /// let batch = decoder.flush().unwrap().unwrap();
581    ///
582    /// // Expect batch containing two columns
583    /// let int32 = batch.column(0).as_primitive::<Int32Type>();
584    /// assert_eq!(int32.values(), &[0, 4]);
585    ///
586    /// let float = batch.column(1).as_primitive::<Float32Type>();
587    /// assert_eq!(float.values(), &[3., 67.53]);
588    /// ```
589    ///
590    /// Or even complex nested types
591    ///
592    /// ```
593    /// # use std::collections::BTreeMap;
594    /// # use std::sync::Arc;
595    /// # use arrow_array::StructArray;
596    /// # use arrow_cast::display::{ArrayFormatter, FormatOptions};
597    /// # use arrow_json::ReaderBuilder;
598    /// # use arrow_schema::{DataType, Field, Fields, Schema};
599    /// # use serde::Serialize;
600    /// #
601    /// #[derive(Serialize)]
602    /// struct MyStruct {
603    ///     int32: i32,
604    ///     list: Vec<f64>,
605    ///     nested: Vec<Option<Nested>>,
606    /// }
607    ///
608    /// impl MyStruct {
609    ///     /// Returns the [`Fields`] for [`MyStruct`]
610    ///     fn fields() -> Fields {
611    ///         let nested = DataType::Struct(Nested::fields());
612    ///         Fields::from([
613    ///             Arc::new(Field::new("int32", DataType::Int32, false)),
614    ///             Arc::new(Field::new_list(
615    ///                 "list",
616    ///                 Field::new("element", DataType::Float64, false),
617    ///                 false,
618    ///             )),
619    ///             Arc::new(Field::new_list(
620    ///                 "nested",
621    ///                 Field::new("element", nested, true),
622    ///                 true,
623    ///             )),
624    ///         ])
625    ///     }
626    /// }
627    ///
628    /// #[derive(Serialize)]
629    /// struct Nested {
630    ///     map: BTreeMap<String, Vec<String>>
631    /// }
632    ///
633    /// impl Nested {
634    ///     /// Returns the [`Fields`] for [`Nested`]
635    ///     fn fields() -> Fields {
636    ///         let element = Field::new("element", DataType::Utf8, false);
637    ///         Fields::from([
638    ///             Arc::new(Field::new_map(
639    ///                 "map",
640    ///                 "entries",
641    ///                 Field::new("key", DataType::Utf8, false),
642    ///                 Field::new_list("value", element, false),
643    ///                 false, // sorted
644    ///                 false, // nullable
645    ///             ))
646    ///         ])
647    ///     }
648    /// }
649    ///
650    /// let data = vec![
651    ///     MyStruct {
652    ///         int32: 34,
653    ///         list: vec![1., 2., 34.],
654    ///         nested: vec![
655    ///             None,
656    ///             Some(Nested {
657    ///                 map: vec![
658    ///                     ("key1".to_string(), vec!["foo".to_string(), "bar".to_string()]),
659    ///                     ("key2".to_string(), vec!["baz".to_string()])
660    ///                 ].into_iter().collect()
661    ///             })
662    ///         ]
663    ///     },
664    ///     MyStruct {
665    ///         int32: 56,
666    ///         list: vec![],
667    ///         nested: vec![]
668    ///     },
669    ///     MyStruct {
670    ///         int32: 24,
671    ///         list: vec![-1., 245.],
672    ///         nested: vec![None]
673    ///     }
674    /// ];
675    ///
676    /// let schema = Schema::new(MyStruct::fields());
677    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
678    /// decoder.serialize(&data).unwrap();
679    /// let batch = decoder.flush().unwrap().unwrap();
680    /// assert_eq!(batch.num_rows(), 3);
681    /// assert_eq!(batch.num_columns(), 3);
682    ///
683    /// // Convert to StructArray to format
684    /// let s = StructArray::from(batch);
685    /// let options = FormatOptions::default().with_null("null");
686    /// let formatter = ArrayFormatter::try_new(&s, &options).unwrap();
687    ///
688    /// assert_eq!(&formatter.value(0).to_string(), "{int32: 34, list: [1.0, 2.0, 34.0], nested: [null, {map: {key1: [foo, bar], key2: [baz]}}]}");
689    /// assert_eq!(&formatter.value(1).to_string(), "{int32: 56, list: [], nested: []}");
690    /// assert_eq!(&formatter.value(2).to_string(), "{int32: 24, list: [-1.0, 245.0], nested: [null]}");
691    /// ```
692    ///
693    /// Note: this ignores any batch size setting, and always decodes all rows
694    ///
695    /// [serde]: https://docs.rs/serde/latest/serde/
696    pub fn serialize<S: Serialize>(&mut self, rows: &[S]) -> Result<(), ArrowError> {
697        self.tape_decoder.serialize(rows)
698    }
699
700    /// True if the decoder is currently part way through decoding a record.
701    pub fn has_partial_record(&self) -> bool {
702        self.tape_decoder.has_partial_row()
703    }
704
705    /// The number of unflushed records, including the partially decoded record (if any).
706    pub fn len(&self) -> usize {
707        self.tape_decoder.num_buffered_rows()
708    }
709
710    /// True if there are no records to flush, i.e. [`Self::len`] is zero.
711    pub fn is_empty(&self) -> bool {
712        self.len() == 0
713    }
714
715    /// Flushes the currently buffered data to a [`RecordBatch`]
716    ///
717    /// Returns `Ok(None)` if no buffered data, i.e. [`Self::is_empty`] is true.
718    ///
719    /// Note: This will return an error if called part way through decoding a record,
720    /// i.e. [`Self::has_partial_record`] is true.
721    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        // First offset is null sentinel
729        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
751/// Decodes a column of JSON values from a [`Tape`] into an [`ArrayRef`]
752///
753/// Implement alongside [`DecoderFactory`] to override or add support for a type.
754pub trait ArrayDecoder: Send {
755    /// Decode elements from `tape` starting at the indexes contained in `pos`
756    ///
757    /// `pos` contains one tape index per output row, so the returned array must have
758    /// exactly `pos.len()` elements and the field's data type. A row's value may be
759    /// [`TapeElement::Null`].
760    fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError>;
761}
762
763/// A trait to create custom decoders for specific data types.
764///
765/// Overrides the reader's decoder for a data type, or adds support for one it does not
766/// handle. The reader-side counterpart of [`EncoderFactory`]; register an
767/// implementation with [`ReaderBuilder::with_decoder_factory`].
768///
769/// # Examples
770///
771/// Decodes `Binary` from a JSON array of integers rather than the default hex string.
772///
773/// ```
774/// use std::sync::Arc;
775/// use arrow_array::{Array, ArrayRef, BinaryArray};
776/// use arrow_array::types::Float64Type;
777/// use arrow_array::cast::AsArray;
778/// use arrow_json::reader::{ArrayDecoder, DecoderContext, DecoderFactory, Tape, TapeElement};
779/// use arrow_json::ReaderBuilder;
780/// use arrow_schema::{ArrowError, DataType, Field, FieldRef, Fields, Schema};
781/// use arrow_schema::extension::EXTENSION_TYPE_NAME_KEY;
782/// use arrow_array::StringArray;
783///
784/// /// Decodes `[104, 105]` into the bytes `b"hi"`
785/// struct IntArrayBinaryDecoder;
786///
787/// impl ArrayDecoder for IntArrayBinaryDecoder {
788///     fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
789///         let mut values: Vec<Option<Vec<u8>>> = Vec::with_capacity(pos.len());
790///         for p in pos {
791///             match tape.get(*p) {
792///                 TapeElement::Null => values.push(None),
793///                 TapeElement::StartList(end) => {
794///                     let mut bytes = Vec::new();
795///                     let mut cur = p + 1;
796///                     while cur < end {
797///                         match tape.get(cur) {
798///                             // JSON text yields `Number`; serde yields `I32`
799///                             TapeElement::Number(idx) => {
800///                                 let s = tape.get_string(idx);
801///                                 bytes.push(s.parse::<u8>().map_err(|e| {
802///                                     ArrowError::JsonError(format!("invalid byte {s}: {e}"))
803///                                 })?);
804///                             }
805///                             TapeElement::I32(v) => bytes.push(v as u8),
806///                             _ => return Err(tape.error(cur, "byte")),
807///                         }
808///                         cur = tape.next(cur, "byte")?;
809///                     }
810///                     values.push(Some(bytes));
811///                 }
812///                 _ => return Err(tape.error(*p, "list of bytes")),
813///             }
814///         }
815///         Ok(Arc::new(BinaryArray::from_iter(values.iter().map(|v| v.as_deref()))))
816///     }
817/// }
818///
819/// /// Upper-cases whatever the reader's own decoder produced
820/// struct ShoutDecoder(Box<dyn ArrayDecoder>);
821///
822/// impl ArrayDecoder for ShoutDecoder {
823///     fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError> {
824///         let inner = self.0.decode(tape, pos)?;
825///         let values = inner.as_string::<i32>();
826///         Ok(Arc::new(StringArray::from_iter(
827///             values.iter().map(|v| v.map(str::to_uppercase)),
828///         )))
829///     }
830/// }
831///
832/// #[derive(Debug)]
833/// struct IntArrayBinaryDecoderFactory;
834///
835/// impl DecoderFactory for IntArrayBinaryDecoderFactory {
836///     fn make_default_decoder(
837///         &self,
838///         ctx: &DecoderContext,
839///         field: &FieldRef,
840///         is_nullable: bool,
841///     ) -> Result<Option<Box<dyn ArrayDecoder>>, ArrowError> {
842///         // Selection can key off metadata, e.g. to recognise an extension type, and
843///         // build on the reader's own decoder for the very same field
844///         if field.metadata().get(EXTENSION_TYPE_NAME_KEY).map(String::as_str)
845///             == Some("apache.shout")
846///         {
847///             let inner = ctx.make_builtin_decoder(field, is_nullable)?;
848///             return Ok(Some(Box::new(ShoutDecoder(inner))));
849///         }
850///
851///         match field.data_type() {
852///             DataType::Binary => Ok(Some(Box::new(IntArrayBinaryDecoder))),
853///             // Returning `None` uses the reader's default decoder
854///             _ => Ok(None),
855///         }
856///     }
857/// }
858///
859/// let nested = Fields::from(vec![Field::new("inner", DataType::Binary, true)]);
860/// let schema = Arc::new(Schema::new(vec![
861///     Field::new("bytes", DataType::Binary, true),
862///     Field::new("float", DataType::Float64, true),
863///     Field::new("nested", DataType::Struct(nested), true),
864///     Field::new("shout", DataType::Utf8, true)
865///         .with_metadata([(EXTENSION_TYPE_NAME_KEY, "apache.shout")]),
866/// ]));
867///
868/// let json = r#"{"bytes": [104, 105], "float": 1.0, "nested": {"inner": [104, 105]}, "shout": "hi"}
869/// {"float": 2.3}
870/// {"bytes": [98], "nested": {"inner": [98]}}
871/// "#;
872///
873/// let batch = ReaderBuilder::new(schema)
874///     .with_decoder_factory(Arc::new(IntArrayBinaryDecoderFactory))
875///     .build(json.as_bytes())
876///     .unwrap()
877///     .next()
878///     .unwrap()
879///     .unwrap();
880///
881/// let bytes = batch.column(0).as_binary::<i32>();
882/// assert_eq!(bytes.value(0), b"hi");
883/// assert!(bytes.is_null(1));
884/// assert_eq!(bytes.value(2), b"b");
885///
886/// // The override applies wherever `Binary` appears, including nested, while types
887/// // the factory declines are decoded as usual
888/// let inner = batch.column(2).as_struct().column(0).as_binary::<i32>();
889/// assert_eq!(inner.value(0), b"hi");
890/// assert_eq!(batch.column(1).as_primitive::<Float64Type>().value(0), 1.0);
891///
892/// // Dispatched on metadata, and decoded by the reader's own `Utf8` decoder
893/// assert_eq!(batch.column(3).as_string::<i32>().value(0), "HI");
894/// ```
895///
896/// [`EncoderFactory`]: crate::EncoderFactory
897pub trait DecoderFactory: std::fmt::Debug + Send + Sync {
898    /// Make a decoder for `field`, or `Ok(None)` to use the reader's default.
899    ///
900    /// Receives the [`FieldRef`] rather than just its [`DataType`] so decoder
901    /// selection can consider the field's metadata, e.g. to identify [extension
902    /// types]. The root of a [`ReaderBuilder::new`] schema is presented as a
903    /// synthesized nameless `Struct` field.
904    ///
905    /// Use [`DecoderContext::make_decoder`] on `ctx` to build child decoders, and
906    /// [`DecoderContext::make_builtin_decoder`] to build on the reader's own decoder
907    /// for this field. Calling `make_decoder` with the field this was invoked with
908    /// recurses back here and loops.
909    ///
910    /// `is_nullable` folds in ancestor nullability, so it may differ from
911    /// `field.is_nullable()` in either direction: a nullable struct widens its
912    /// children, while a run-end encoded array narrows its values.
913    ///
914    /// [extension types]: https://arrow.apache.org/docs/format/Columnar.html#extension-types
915    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
925/// Validates the output of a [`DecoderFactory`] decoder before it reaches the arrays
926/// built from it, some of which are constructed without further checks.
927struct 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
953/// Context for decoder creation, containing configuration.
954///
955/// This context is passed through the decoder creation process and contains
956/// all the configuration needed to create decoders recursively.
957pub struct DecoderContext {
958    /// Whether to coerce primitives to strings
959    coerce_primitive: bool,
960    /// Whether to validate struct fields strictly
961    strict_mode: bool,
962    /// How to decode struct fields
963    struct_mode: StructMode,
964    /// Whether to treat columns with incompatible types as missing (i.e. NULL)
965    ignore_type_conflicts: bool,
966    /// An optional hook for customizing decoding behavior
967    decoder_factory: Option<Arc<dyn DecoderFactory>>,
968}
969
970impl DecoderContext {
971    /// Returns whether to coerce primitive types (e.g., number to string)
972    pub fn coerce_primitive(&self) -> bool {
973        self.coerce_primitive
974    }
975
976    /// Returns whether to validate struct fields strictly
977    pub fn strict_mode(&self) -> bool {
978        self.strict_mode
979    }
980
981    /// Returns how to decode struct fields
982    pub fn struct_mode(&self) -> StructMode {
983        self.struct_mode
984    }
985
986    /// Returns whether to treat columns with incompatible types as missing (i.e. NULL)
987    pub fn ignore_type_conflicts(&self) -> bool {
988        self.ignore_type_conflicts
989    }
990
991    /// Returns the optional hook for customizing decoding behavior
992    pub fn decoder_factory(&self) -> Option<&Arc<dyn DecoderFactory>> {
993        self.decoder_factory.as_ref()
994    }
995
996    /// Create a decoder for a type.
997    ///
998    /// This is the standard way to create child decoders from within a decoder
999    /// implementation.
1000    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    /// Create the decoder the reader would use for `field`, ignoring any
1009    /// [`DecoderFactory`].
1010    ///
1011    /// Lets a factory build on the reader's own decoder, including for the field it was
1012    /// asked about — which [`Self::make_decoder`] cannot do without looping.
1013    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        // Test with different batch sizes to test for boundary conditions
1172        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            // Test with different buffer sizes to test for boundary conditions
1186            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        // The JSON input contains field "a" with different string lengths.
1307        // According to the implementation in the decoder:
1308        // - For a string, capacity is only increased if its length > 12 bytes.
1309        // Therefore, for:
1310        // Row 1: "short" (5 bytes) -> capacity += 0
1311        // Row 2: "this is definitely long" (24 bytes) -> capacity += 24
1312        // Row 3: "hello" (5 bytes) -> capacity += 0
1313        // Row 4: "\nfoobar😀asfgÿ" (17 bytes) -> capacity += 17
1314        // Expected total capacity = 24 + 17 = 41
1315        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        // Get the first column ("a") as a StringViewArray.
1333        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        // Retrieve the underlying data buffer from the array.
1340        // The builder pre-allocates capacity based on the sum of lengths for long strings.
1341        let data_buffer = string_view_array.to_data().buffers()[0].len();
1342
1343        // Check that the allocated capacity is at least what we expected.
1344        // (The actual buffer may be larger than expected due to rounding or internal allocation strategies.)
1345        assert!(
1346            data_buffer >= expected_capacity,
1347            "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1348        );
1349
1350        // Additionally, verify that the decoded values are correct.
1351        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 the memory capacity allocation logic when converting numeric types to strings.
1358    #[test]
1359    fn test_numeric_view_allocation() {
1360        // For numeric types, the expected capacity calculation is as follows:
1361        // Row 1: 123456789  -> Number converts to the string "123456789" (length 9), 9 <= 12, so no capacity is added.
1362        // Row 2: 1000000000000 -> Treated as an I64 number; its string is "1000000000000" (length 13),
1363        //                        which is >12 and its absolute value is > 999_999_999_999, so 13 bytes are added.
1364        // Row 3: 3.1415 -> F32 number, a fixed estimate of 10 bytes is added.
1365        // Row 4: 2.718281828459045 -> F64 number, a fixed estimate of 10 bytes is added.
1366        // Total expected capacity = 13 + 10 + 10 = 33 bytes.
1367        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        // Check that the underlying data buffer capacity is at least the expected value.
1388        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        // Verify that the converted string values are correct.
1395        // Note: The format of the number converted to a string should match the actual implementation.
1396        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        // Rounding half away from zero, exponent notation (a valid JSON
1774        // number syntax), surrounding whitespace in strings and negative
1775        // scales are all accepted
1776        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        // Invalid and out-of-range values are errors (or nulls when type
1812        // conflicts are ignored), never panics
1813        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)] // Takes too long
1953    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            // Test correctly enforced nullability
2112            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()); // Should error as not nullable
2122
2123            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()); // Should error as not nullable
2128
2129            // Test nulls in nullable parent can mask nulls in non-nullable child
2130            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            // Non-nullable child can still contain null as masked by parent
2144            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        // build expected output
2577        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        // compare `a` with result from json reader
2620        let batch = reader.next().unwrap().unwrap();
2621        let read = batch.column(0);
2622        assert_eq!(read.len(), 6);
2623        // compare the arrays the long way around, to better detect differences
2624        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        // compare list null buffers
2628        assert_eq!(read.nulls(), expected.nulls());
2629        // build struct from list
2630        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        // test struct's nulls
2638        assert_eq!(struct_array.nulls(), expected_struct_array.nulls());
2639        // test struct's fields
2640        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        // Check offsets and sizes
2676        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        // Row 0: [1, 2, 3]
2688        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        // Row 1: [4, null]
2694        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        // Row 2: null
2702        assert!(list_view.is_null(2));
2703
2704        // Row 3: [6]
2705        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        // Row 4: []
2711        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)); // string -> null
2912        assert!(col.is_null(2)); // number -> null
2913        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        // read record batches
3314        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        // Coercing an f32 into a string column must render the value, not its raw bit pattern.
3340        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    // Parse the given `row` in `struct_mode` as a type given by fields.
3353    //
3354    // If as_struct == true, wrap the fields in a Struct field with name "r".
3355    // If as_struct == false, wrap the fields in a Schema.
3356    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 cases:
3519    // [] -> RecordBatch row with no entries.  Schema = [('a', Int32)] -> Error
3520    // [] -> RecordBatch row with no entries. Schema = [('r', [('a', Int32)])] -> Error
3521    // [] -> StructArray row with no entries. Fields [('a', Int32')] -> Error
3522    // [[]] -> RecordBatch row with empty struct entry. Schema = [('r', [('a', Int32)])] -> Error
3523    #[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        // Missing fields
3542        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        // Wrong values
3587        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, // not nullable
3642                    )),
3643                    false, // not sorted
3644                ),
3645                true, // nullable
3646            ),
3647            Field::new(
3648                "struct",
3649                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3650                true,
3651            ),
3652        ]);
3653
3654        // A compatible value for each schema field above, in schema order
3655        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        // Create a set of JSON rows that rotates each value past every field
3669        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        // NOTE: NullArray doesn't materialize any values (they're all NULL by definition)
3691        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                // Both map and struct can parse
3831                &[
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                // Both map and struct can parse
3854                &[
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, // not nullable
3900                    )),
3901                    false, // not sorted
3902                ),
3903                false, // not nullable
3904            ),
3905            Field::new(
3906                "struct",
3907                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3908                false,
3909            ),
3910        ];
3911
3912        // Every field above will have a type conflict with at least one of these values
3913        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, // not nullable
3961                    )),
3962                    false, // not sorted
3963                ),
3964                true, // not nullable
3965            ),
3966            Field::new(
3967                "struct",
3968                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3969                true,
3970            ),
3971        ];
3972
3973        // Every field above will have a type conflict with at least one of these values
3974        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        // 5 logical values compressed into 2 runs
4009        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        // 5 logical values: "x", null, null, null, "y" → 3 runs
4040        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        // No compression: 3 unique values → 3 runs
4109        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        // The outer REE compresses whole rows, while the inner REE compresses the
4145        // repeated string values produced by decoding those rows.
4146        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        // Three logical rows compress to two outer runs: ["x", "x"] and ["y"].
4161        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        // The physical values of the outer REE are themselves a two-element REE.
4166        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    /// Declining everything must be indistinguishable from no factory at all.
4199    #[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        // Consulted for the root struct and every leaf, with names attached; the
4244        // root is the documented synthesized nameless struct
4245        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    /// A custom decoder's output is not trusted: the arrays built from it are
4269    /// constructed without further validation, so a mismatch must be caught here.
4270    #[test]
4271    fn test_decoder_factory_output_is_validated() {
4272        /// Returns either the wrong data type or the wrong number of rows
4273        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}