From 9ed71fcdc675582b7a3fdd7c20f7cfb59fdcf102 Mon Sep 17 00:00:00 2001 From: Liam Bao Date: Fri, 3 Apr 2026 17:37:01 -0400 Subject: [PATCH] [Json] Support FixedSizeList in json decoder --- arrow-json/src/lib.rs | 38 +++++++ arrow-json/src/reader/list_array.rs | 85 ++++++++++++++- arrow-json/src/reader/mod.rs | 161 +++++++++++++++++++++++++++- 3 files changed, 282 insertions(+), 2 deletions(-) diff --git a/arrow-json/src/lib.rs b/arrow-json/src/lib.rs index 201c3cd80749..9f2a9e3a81ca 100644 --- a/arrow-json/src/lib.rs +++ b/arrow-json/src/lib.rs @@ -390,4 +390,42 @@ mod tests { assert_list_view_roundtrip::(); assert_list_view_roundtrip::(); } + + #[test] + fn test_json_roundtrip_fixed_size_list() { + let inner = Arc::new(Field::new("item", DataType::Int32, true)); + let schema = Arc::new(Schema::new(vec![ + Field::new("flat", DataType::FixedSizeList(inner.clone(), 3), true), + Field::new( + "nested", + DataType::FixedSizeList( + Arc::new(Field::new("item", DataType::FixedSizeList(inner, 2), true)), + 2, + ), + true, + ), + ])); + + let input = r#"{"flat":[1,2,3],"nested":[[1,2],[3,4]]} +{"flat":[4,null,5]} +{"flat":[6,7,8],"nested":[[null,5],[6,null]]} +"# + .as_bytes(); + + let batches: Vec = ReaderBuilder::new(schema.clone()) + .with_batch_size(1024) + .build(Cursor::new(input)) + .unwrap() + .collect::, _>>() + .unwrap(); + + let mut output = Vec::new(); + let mut writer = WriterBuilder::new().build::<_, LineDelimited>(&mut output); + for batch in &batches { + writer.write(batch).unwrap(); + } + writer.finish().unwrap(); + + assert_eq!(input, &output); + } } diff --git a/arrow-json/src/reader/list_array.rs b/arrow-json/src/reader/list_array.rs index 113e628541c6..6c56dbad8a4b 100644 --- a/arrow-json/src/reader/list_array.rs +++ b/arrow-json/src/reader/list_array.rs @@ -19,7 +19,7 @@ use std::marker::PhantomData; use std::sync::Arc; use arrow_array::builder::BooleanBufferBuilder; -use arrow_array::{ArrayRef, GenericListArray, OffsetSizeTrait, make_array}; +use arrow_array::{ArrayRef, FixedSizeListArray, GenericListArray, OffsetSizeTrait, make_array}; use arrow_buffer::buffer::NullBuffer; use arrow_buffer::{Buffer, OffsetBuffer, ScalarBuffer}; use arrow_data::ArrayDataBuilder; @@ -140,3 +140,86 @@ impl ArrayDecoder for ListLikeArrayDeco } } } + +pub struct FixedSizeListArrayDecoder { + field: FieldRef, + size: i32, + decoder: Box, + ignore_type_conflicts: bool, + is_nullable: bool, +} + +impl FixedSizeListArrayDecoder { + pub fn new( + ctx: &DecoderContext, + data_type: &DataType, + is_nullable: bool, + ) -> Result { + let (field, size) = match data_type { + DataType::FixedSizeList(f, s) => (f, *s), + _ => unreachable!(), + }; + let decoder = ctx.make_decoder(field.data_type(), field.is_nullable())?; + + Ok(Self { + field: field.clone(), + size, + decoder, + ignore_type_conflicts: ctx.ignore_type_conflicts(), + is_nullable, + }) + } +} + +impl ArrayDecoder for FixedSizeListArrayDecoder { + fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result { + let expected = self.size as usize; + let mut child_pos = Vec::with_capacity(pos.len() * expected); + + let mut nulls = self + .is_nullable + .then(|| BooleanBufferBuilder::new(pos.len())); + + for p in pos { + let end_idx = match (tape.get(*p), nulls.as_mut()) { + (TapeElement::StartList(end_idx), None) => end_idx, + (TapeElement::StartList(end_idx), Some(nulls)) => { + nulls.append(true); + end_idx + } + (TapeElement::Null, Some(nulls)) => { + nulls.append(false); + child_pos.resize(child_pos.len() + expected, 0); + continue; + } + (_, Some(nulls)) if self.ignore_type_conflicts => { + nulls.append(false); + child_pos.resize(child_pos.len() + expected, 0); + continue; + } + _ => return Err(tape.error(*p, "[")), + }; + + let child_start = child_pos.len(); + let mut cur_idx = *p + 1; + while cur_idx < end_idx { + child_pos.push(cur_idx); + cur_idx = tape.next(cur_idx, "fixed-size list value")?; + } + + let actual = child_pos.len() - child_start; + if actual != expected { + return Err(ArrowError::JsonError(format!( + "Incorrect number of elements for FixedSizeList, \ + expected {expected} but got {actual}" + ))); + } + } + + let values = self.decoder.decode(tape, &child_pos)?; + let nulls = nulls.as_mut().map(|x| NullBuffer::new(x.finish())); + + let array = FixedSizeListArray::try_new(self.field.clone(), self.size, values, nulls)?; + Ok(Arc::new(array)) + } +} diff --git a/arrow-json/src/reader/mod.rs b/arrow-json/src/reader/mod.rs index 62c13c70ed99..0e04ef4c1074 100644 --- a/arrow-json/src/reader/mod.rs +++ b/arrow-json/src/reader/mod.rs @@ -151,7 +151,9 @@ use crate::reader::binary_array::{ }; use crate::reader::boolean_array::BooleanArrayDecoder; use crate::reader::decimal_array::DecimalArrayDecoder; -use crate::reader::list_array::{ListArrayDecoder, ListViewArrayDecoder}; +use crate::reader::list_array::{ + FixedSizeListArrayDecoder, ListArrayDecoder, ListViewArrayDecoder, +}; use crate::reader::map_array::MapArrayDecoder; use crate::reader::null_array::NullArrayDecoder; use crate::reader::primitive_array::PrimitiveArrayDecoder; @@ -835,6 +837,7 @@ fn make_decoder( DataType::LargeList(_) => Ok(Box::new(ListArrayDecoder::::new(ctx, data_type, is_nullable)?)), DataType::ListView(_) => Ok(Box::new(ListViewArrayDecoder::::new(ctx, data_type, is_nullable)?)), DataType::LargeListView(_) => Ok(Box::new(ListViewArrayDecoder::::new(ctx, data_type, is_nullable)?)), + DataType::FixedSizeList(_, _) => Ok(Box::new(FixedSizeListArrayDecoder::new(ctx, data_type, is_nullable)?)), DataType::Struct(_) => Ok(Box::new(StructArrayDecoder::new(ctx, data_type, is_nullable)?)), DataType::Binary => Ok(Box::new(BinaryArrayDecoder::::default())), DataType::LargeBinary => Ok(Box::new(BinaryArrayDecoder::::default())), @@ -2308,6 +2311,152 @@ mod tests { assert_read_list_view::(); } + #[test] + fn test_fixed_size_list() { + let buf = r#" + {"a": [1, 2, 3]} + {"a": [4, 5, 6]} + {"a": [7, 8, 9]} + "#; + + let field = Field::new_list_field(DataType::Int32, true); + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::FixedSizeList(Arc::new(field), 3), + false, + )])); + + let batches = do_read(buf, 1024, false, false, schema); + assert_eq!(batches.len(), 1); + + let col = batches[0].column(0).as_fixed_size_list(); + assert_eq!(col.len(), 3); + assert_eq!(col.value_length(), 3); + + let values = col.values().as_primitive::(); + assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8, 9]); + } + + #[test] + fn test_fixed_size_list_nullable() { + let buf = r#" + {"a": [1, 2]} + {"a": null} + {"a": [3, null]} + "#; + + let field = Field::new_list_field(DataType::Int32, true); + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::FixedSizeList(Arc::new(field), 2), + true, + )])); + + let batches = do_read(buf, 1024, false, false, schema); + assert_eq!(batches.len(), 1); + + let col = batches[0].column(0).as_fixed_size_list(); + assert_eq!(col.len(), 3); + assert!(col.is_valid(0)); + assert!(col.is_null(1)); + assert!(col.is_valid(2)); + + let values = col.values().as_primitive::(); + assert_eq!(values.value(0), 1); + assert_eq!(values.value(1), 2); + assert_eq!(values.value(4), 3); + assert!(values.is_null(5)); + } + + #[test] + fn test_fixed_size_list_wrong_size() { + let buf = r#"{"a": [1, 2, 3]}"#; + + let field = Field::new_list_field(DataType::Int32, true); + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::FixedSizeList(Arc::new(field), 2), + false, + )])); + + let err = ReaderBuilder::new(schema) + .build(Cursor::new(buf.as_bytes())) + .unwrap() + .next() + .unwrap() + .unwrap_err(); + + assert!(err.to_string().contains("expected 2 but got 3"), "{}", err); + } + + #[test] + fn test_fixed_size_list_nested() { + let buf = r#" + {"a": [[1, 2], [3, 4]]} + {"a": [[5, 6], [7, 8]]} + "#; + + let inner_field = Field::new_list_field(DataType::Int32, true); + let inner_type = DataType::FixedSizeList(Arc::new(inner_field), 2); + let outer_field = Arc::new(Field::new_list_field(inner_type.clone(), true)); + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::FixedSizeList(outer_field, 2), + false, + )])); + + let batches = do_read(buf, 1024, false, false, schema); + assert_eq!(batches.len(), 1); + + let col = batches[0].column(0).as_fixed_size_list(); + assert_eq!(col.len(), 2); + assert_eq!(col.value_length(), 2); + + let inner = col.values().as_fixed_size_list(); + assert_eq!(inner.len(), 4); + assert_eq!(inner.value_length(), 2); + + let values = inner.values().as_primitive::(); + assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8]); + } + + #[test] + fn test_fixed_size_list_ignore_type_conflicts() { + let field = Field::new("item", DataType::Int32, true); + let schema = Arc::new(Schema::new(vec![Field::new( + "a", + DataType::FixedSizeList(Arc::new(field), 2), + true, + )])); + + let json = vec![ + json!({"a": [1, 2]}), + json!({"a": "not a list"}), + json!({"a": 42}), + json!({"a": [6, 7]}), + ]; + + let mut decoder = ReaderBuilder::new(schema) + .with_ignore_type_conflicts(true) + .build_decoder() + .unwrap(); + decoder.serialize(&json).unwrap(); + let batch = decoder.flush().unwrap().unwrap(); + + let col = batch.column(0).as_fixed_size_list(); + assert_eq!(col.len(), 4); + assert!(col.is_valid(0)); + assert!(col.is_null(1)); // string -> null + assert!(col.is_null(2)); // number -> null + assert!(col.is_valid(3)); + + let values = col.values().as_primitive::(); + assert_eq!(values.value(0), 1); + assert_eq!(values.value(1), 2); + assert_eq!(values.value(6), 6); + assert_eq!(values.value(7), 7); + } + #[test] fn test_skip_empty_lines() { let schema = Schema::new(vec![Field::new("a", DataType::Int64, true)]); @@ -3256,6 +3405,11 @@ mod tests { DataType::List(Arc::new(Field::new("item", DataType::Int32, true))), false, ), + Field::new( + "fixed_size_list", + DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2), + false, + ), Field::new( "map", DataType::Map( @@ -3312,6 +3466,11 @@ mod tests { DataType::List(Arc::new(Field::new("item", DataType::Int32, true))), true, ), + Field::new( + "fixed_size_list", + DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2), + true, + ), Field::new( "map", DataType::Map(