Skip to main content

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