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