This is an automated email from the ASF dual-hosted git repository.

guan404ming pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/mahout.git


The following commit(s) were added to refs/heads/main by this push:
     new 91335ed19 [Bug][QDP] Fix List<T> readers failing on nullable outer 
rows (#1402)
91335ed19 is described below

commit 91335ed1939a1b8f6c5eb5d7451d2d21e77da8ab
Author: ChenChen Lai <[email protected]>
AuthorDate: Thu Jun 18 17:51:29 2026 +0800

    [Bug][QDP] Fix List<T> readers failing on nullable outer rows (#1402)
    
    * Fix List<T> readers failing on nullable outer rows
    
    * Fix ParquetStreamingReader rejecting all-null leading batches
---
 qdp/qdp-core/src/readers/arrow_ipc.rs | 215 ++++++++++++++++++++--
 qdp/qdp-core/src/readers/parquet.rs   | 328 +++++++++++++++++++++++++++++++---
 2 files changed, 510 insertions(+), 33 deletions(-)

diff --git a/qdp/qdp-core/src/readers/arrow_ipc.rs 
b/qdp/qdp-core/src/readers/arrow_ipc.rs
index 39d9d9148..94ea1faaa 100644
--- a/qdp/qdp-core/src/readers/arrow_ipc.rs
+++ b/qdp/qdp-core/src/readers/arrow_ipc.rs
@@ -150,16 +150,12 @@ impl DataReader for ArrowIPCReader {
                             MahoutError::Io("Failed to downcast to 
ListArray".to_string())
                         })?;
 
+                    // Phase 1: find sample_size from non-null rows and 
validate consistency.
                     for i in 0..list_array.len() {
-                        let value_array = list_array.value(i);
-                        let float_array = value_array
-                            .as_any()
-                            .downcast_ref::<Float64Array>()
-                            .ok_or_else(|| {
-                                MahoutError::Io("List values must be 
Float64".to_string())
-                            })?;
-
-                        let current_size = float_array.len();
+                        if list_array.is_null(i) {
+                            continue;
+                        }
+                        let current_size = list_array.value_length(i) as usize;
 
                         if let Some(expected) = sample_size {
                             if current_size != expected {
@@ -180,10 +176,53 @@ impl DataReader for ArrowIPCReader {
                                 })?;
                             all_data.reserve(new_capacity);
                         }
+                    }
 
+                    // Phase 2: collect data, handling null outer rows per 
NullHandling policy.
+                    if list_array.null_count() == 0 {
+                        let values = list_array.values();
+                        let float_array = values
+                            .as_any()
+                            .downcast_ref::<Float64Array>()
+                            .ok_or_else(|| MahoutError::Io("Values must be 
Float64".to_string()))?;
                         handle_float64_nulls(&mut all_data, float_array, 
self.null_handling)?;
-
-                        num_samples += 1;
+                        num_samples += list_array.len();
+                    } else {
+                        for i in 0..list_array.len() {
+                            if list_array.is_null(i) {
+                                match self.null_handling {
+                                    NullHandling::Reject => {
+                                        return Err(MahoutError::InvalidInput(
+                                            "Null outer row in List column. 
Use \
+                                             NullHandling::FillZero to replace 
with zeros, \
+                                             or clean the data at the source."
+                                                .to_string(),
+                                        ));
+                                    }
+                                    NullHandling::FillZero => {
+                                        if let Some(ss) = sample_size {
+                                            
all_data.extend(std::iter::repeat_n(0.0_f64, ss));
+                                            num_samples += 1;
+                                        }
+                                        // sample_size unknown: skip this null 
row without counting it.
+                                    }
+                                }
+                            } else {
+                                let value_array = list_array.value(i);
+                                let float_array = value_array
+                                    .as_any()
+                                    .downcast_ref::<Float64Array>()
+                                    .ok_or_else(|| {
+                                        MahoutError::Io("List values must be 
Float64".to_string())
+                                    })?;
+                                handle_float64_nulls(
+                                    &mut all_data,
+                                    float_array,
+                                    self.null_handling,
+                                )?;
+                                num_samples += 1;
+                            }
+                        }
                     }
                 }
 
@@ -202,3 +241,157 @@ impl DataReader for ArrowIPCReader {
         Ok((all_data, num_samples, sample_size))
     }
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use arrow::array::{ArrayRef, Float64Builder, ListBuilder, RecordBatch};
+    use arrow::datatypes::{DataType, Field, Schema};
+    use arrow::ipc::writer::FileWriter as ArrowIpcFileWriter;
+    use std::fs;
+    use std::sync::Arc;
+    use std::sync::atomic::{AtomicUsize, Ordering};
+
+    static TEST_FILE_COUNTER: AtomicUsize = AtomicUsize::new(0);
+
+    struct TempTestFile {
+        path: std::path::PathBuf,
+    }
+
+    impl TempTestFile {
+        fn new() -> Self {
+            let count = TEST_FILE_COUNTER.fetch_add(1, Ordering::SeqCst);
+            let path = std::env::temp_dir().join(format!(
+                "mahout_test_arrow_ipc_{}_{}.arrow",
+                std::process::id(),
+                count
+            ));
+            Self { path }
+        }
+
+        fn path(&self) -> &std::path::Path {
+            &self.path
+        }
+    }
+
+    impl Drop for TempTestFile {
+        fn drop(&mut self) {
+            let _ = fs::remove_file(&self.path);
+        }
+    }
+
+    fn write_test_arrow_ipc(schema: Arc<Schema>, arrays: Vec<ArrayRef>) -> 
TempTestFile {
+        let file = TempTestFile::new();
+        let batch = RecordBatch::try_new(schema.clone(), arrays).unwrap();
+        let os_file = fs::File::create(file.path()).unwrap();
+        let mut writer = ArrowIpcFileWriter::try_new(os_file, 
&schema).unwrap();
+        writer.write(&batch).unwrap();
+        writer.finish().unwrap();
+        file
+    }
+
+    fn write_ipc_list_with_null_outer_middle() -> TempTestFile {
+        // [[1.0, 2.0], null, [3.0, 4.0]]
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.append(false); // null outer row
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        write_test_arrow_ipc(schema, vec![array])
+    }
+
+    fn write_ipc_list_with_null_outer_first() -> TempTestFile {
+        // [null, [1.0, 2.0], [3.0, 4.0]] — null row at position 0
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.append(false); // null outer row at position 0
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        write_test_arrow_ipc(schema, vec![array])
+    }
+
+    #[test]
+    fn test_arrow_ipc_reader_null_outer_row_middle_fill_zero() {
+        // [[1,2], null, [3,4]] with FillZero → [1,2, 0,0, 3,4]
+        let file = write_ipc_list_with_null_outer_middle();
+        let mut reader = ArrowIPCReader::new(file.path(), 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        assert_eq!(data, vec![1.0, 2.0, 0.0, 0.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 3);
+        assert_eq!(sample_size, 2);
+    }
+
+    #[test]
+    fn test_arrow_ipc_reader_null_outer_row_middle_reject() {
+        // [[1,2], null, [3,4]] with Reject → error
+        let file = write_ipc_list_with_null_outer_middle();
+        let mut reader = ArrowIPCReader::new(file.path(), 
NullHandling::Reject).unwrap();
+        let result = reader.read_batch();
+        assert!(result.is_err());
+        assert!(result.unwrap_err().to_string().contains("Null outer row"));
+    }
+
+    #[test]
+    fn test_arrow_ipc_reader_null_outer_row_first_fill_zero() {
+        // [null, [1,2], [3,4]] — null at row 0 must not corrupt sample_size
+        let file = write_ipc_list_with_null_outer_first();
+        let mut reader = ArrowIPCReader::new(file.path(), 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        assert_eq!(data, vec![0.0, 0.0, 1.0, 2.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 3);
+        assert_eq!(sample_size, 2);
+    }
+
+    #[test]
+    fn test_arrow_ipc_reader_cross_batch_all_null_first_fill_zero() {
+        // Batch 1: [null, null]   — sample_size unknown
+        // Batch 2: [[1,2], [3,4]] — sample_size established here
+        // All-null leading batch must not corrupt num_samples.
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut b1 = ListBuilder::new(Float64Builder::new());
+        b1.append(false);
+        b1.append(false);
+        let batch1 =
+            RecordBatch::try_new(schema.clone(), vec![Arc::new(b1.finish()) as 
ArrayRef]).unwrap();
+
+        let mut b2 = ListBuilder::new(Float64Builder::new());
+        b2.values().append_slice(&[1.0, 2.0]);
+        b2.append(true);
+        b2.values().append_slice(&[3.0, 4.0]);
+        b2.append(true);
+        let batch2 =
+            RecordBatch::try_new(schema.clone(), vec![Arc::new(b2.finish()) as 
ArrayRef]).unwrap();
+
+        let file = TempTestFile::new();
+        {
+            let os_file = fs::File::create(file.path()).unwrap();
+            let mut writer = ArrowIpcFileWriter::try_new(os_file, 
&schema).unwrap();
+            writer.write(&batch1).unwrap();
+            writer.write(&batch2).unwrap();
+            writer.finish().unwrap();
+        }
+
+        let mut reader = ArrowIPCReader::new(file.path(), 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        assert_eq!(data, vec![1.0, 2.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 2);
+        assert_eq!(sample_size, 2);
+    }
+}
diff --git a/qdp/qdp-core/src/readers/parquet.rs 
b/qdp/qdp-core/src/readers/parquet.rs
index 30f2deead..ac7872213 100644
--- a/qdp/qdp-core/src/readers/parquet.rs
+++ b/qdp/qdp-core/src/readers/parquet.rs
@@ -290,8 +290,12 @@ impl<T: FloatElem> DataReader<T> for ParquetReader<T> {
                             MahoutError::Io("Failed to downcast to 
ListArray".to_string())
                         })?;
 
-                    // Validate all rows have a consistent sample size.
+                    // Validate non-null rows have a consistent sample size.
+                    // Null outer rows return value_length 0, so they must be 
skipped here.
                     for i in 0..list_array.len() {
+                        if list_array.is_null(i) {
+                            continue;
+                        }
                         let row_len = list_array.value_length(i) as usize;
                         if let Some(expected) = sample_size {
                             if row_len != expected {
@@ -306,15 +310,50 @@ impl<T: FloatElem> DataReader<T> for ParquetReader<T> {
                         }
                     }
 
-                    // Cast the entire flat buffer once (avoids N per-row 
allocations
-                    // on cross-dtype reads) then extend all_data in one pass.
-                    let flat = list_flat_values(list_array);
-                    extend_floats::<<T as ArrowPrimitive>::ArrowType>(
-                        &mut all_data,
-                        &*flat,
-                        self.null_handling,
-                    )?;
-                    num_samples += list_array.len();
+                    if list_array.null_count() == 0 {
+                        // Fast path: no null outer rows; use flat buffer.
+                        let flat = list_flat_values(list_array);
+                        extend_floats::<<T as ArrowPrimitive>::ArrowType>(
+                            &mut all_data,
+                            &*flat,
+                            self.null_handling,
+                        )?;
+                        num_samples += list_array.len();
+                    } else {
+                        // Null outer rows present; handle per NullHandling 
policy.
+                        // If sample_size is still unknown (every row in this 
batch is null),
+                        // FillZero cannot determine how many zeros to write — 
those null rows
+                        // are skipped and not counted in num_samples.
+                        for i in 0..list_array.len() {
+                            if list_array.is_null(i) {
+                                match self.null_handling {
+                                    NullHandling::Reject => {
+                                        return Err(MahoutError::InvalidInput(
+                                            "Null outer row in List column. 
Use \
+                                             NullHandling::FillZero to replace 
with zeros, \
+                                             or clean the data at the source."
+                                                .to_string(),
+                                        ));
+                                    }
+                                    NullHandling::FillZero => {
+                                        if let Some(ss) = sample_size {
+                                            
all_data.extend(std::iter::repeat_n(T::default(), ss));
+                                            num_samples += 1;
+                                        }
+                                        // sample_size unknown: skip this null 
row.
+                                    }
+                                }
+                            } else {
+                                let row = list_array.value(i);
+                                extend_floats::<<T as 
ArrowPrimitive>::ArrowType>(
+                                    &mut all_data,
+                                    &*row,
+                                    self.null_handling,
+                                )?;
+                                num_samples += 1;
+                            }
+                        }
+                    }
                 }
                 DataType::FixedSizeList(_, size) => {
                     let list_array = column
@@ -487,7 +526,7 @@ impl<T: FloatElem> DataReader<T> for 
ParquetStreamingReader<T> {
                 break;
             }
             all_data.extend_from_slice(&buffer[..written]);
-            num_samples += written / self.sample_size.unwrap_or(1);
+            num_samples += written / self.sample_size.unwrap_or(1).max(1);
         }
 
         let sample_size = self
@@ -561,10 +600,35 @@ impl<T: FloatElem> StreamingDataReader<T> for 
ParquetStreamingReader<T> {
                                 continue;
                             }
 
-                            let current_sample_size = 
list_array.value_length(0) as usize;
-
-                            // Validate all rows in this batch have a 
consistent sample size.
-                            for i in 1..list_array.len() {
+                            // Find sample_size from the first non-null row.
+                            // Null outer rows return value_length 0 and must 
be skipped.
+                            let first_non_null =
+                                (0..list_array.len()).find(|&i| 
!list_array.is_null(i));
+                            let current_sample_size = match first_non_null {
+                                Some(i) => list_array.value_length(i) as usize,
+                                None => match self.sample_size {
+                                    // All rows null but sample_size known 
from an earlier batch.
+                                    Some(ss) => ss,
+                                    // All rows null and sample_size unknown.
+                                    None => {
+                                        if self.null_handling == 
NullHandling::Reject {
+                                            return 
Err(MahoutError::InvalidInput(
+                                                "Null outer row in List 
column. Use \
+                                                 NullHandling::FillZero to 
replace with \
+                                                 zeros, or clean the data at 
the source."
+                                                    .to_string(),
+                                            ));
+                                        }
+                                        continue;
+                                    }
+                                },
+                            };
+
+                            // Validate all non-null rows in this batch.
+                            for i in 0..list_array.len() {
+                                if list_array.is_null(i) {
+                                    continue;
+                                }
                                 let row_len = list_array.value_length(i) as 
usize;
                                 if row_len != current_sample_size {
                                     return 
Err(MahoutError::InvalidInput(format!(
@@ -574,13 +638,45 @@ impl<T: FloatElem> StreamingDataReader<T> for 
ParquetStreamingReader<T> {
                                 }
                             }
 
-                            // Cast the entire flat buffer once (avoids N 
per-row allocations
-                            // on cross-dtype reads).
-                            let flat = list_flat_values(list_array);
-                            let batch_values = collect_floats::<<T as 
ArrowPrimitive>::ArrowType>(
-                                &*flat,
-                                self.null_handling,
-                            )?;
+                            let batch_values = if list_array.null_count() == 0 
{
+                                // Fast path: no null outer rows; use flat 
buffer.
+                                let flat = list_flat_values(list_array);
+                                collect_floats::<<T as 
ArrowPrimitive>::ArrowType>(
+                                    &*flat,
+                                    self.null_handling,
+                                )?
+                            } else {
+                                // Null outer rows present; handle per 
NullHandling policy.
+                                let mut vals = Vec::new();
+                                for i in 0..list_array.len() {
+                                    if list_array.is_null(i) {
+                                        match self.null_handling {
+                                            NullHandling::Reject => {
+                                                return 
Err(MahoutError::InvalidInput(
+                                                    "Null outer row in List 
column. Use \
+                                                     NullHandling::FillZero to 
replace with \
+                                                     zeros, or clean the data 
at the source."
+                                                        .to_string(),
+                                                ));
+                                            }
+                                            NullHandling::FillZero => {
+                                                
vals.extend(std::iter::repeat_n(
+                                                    T::default(),
+                                                    current_sample_size,
+                                                ));
+                                            }
+                                        }
+                                    } else {
+                                        let row = list_array.value(i);
+                                        extend_floats::<<T as 
ArrowPrimitive>::ArrowType>(
+                                            &mut vals,
+                                            &*row,
+                                            self.null_handling,
+                                        )?;
+                                    }
+                                }
+                                vals
+                            };
 
                             (current_sample_size, batch_values)
                         }
@@ -969,6 +1065,40 @@ mod tests {
 
     // --- NullHandling tests ---
 
+    fn write_list_parquet_with_null_outer_middle() -> TempTestFile {
+        // [[1.0, 2.0], null, [3.0, 4.0]]
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.append(false); // null outer row
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        write_test_parquet(schema, vec![array])
+    }
+
+    fn write_list_parquet_with_null_outer_first() -> TempTestFile {
+        // [null, [1.0, 2.0], [3.0, 4.0]] — null row at position 0 seeds 
sample_size
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.append(false); // null outer row at position 0
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        write_test_parquet(schema, vec![array])
+    }
+
     fn write_list_parquet_with_nulls() -> TempTestFile {
         let item_field = Arc::new(Field::new("item", DataType::Float64, true));
         let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
@@ -1119,4 +1249,158 @@ mod tests {
         };
         assert!(err_msg.contains("Expected List<Float32> or List<Float64>"));
     }
+
+    // --- Null outer row tests ---
+
+    #[test]
+    fn test_parquet_reader_null_outer_row_middle_fill_zero() {
+        // [[1,2], null, [3,4]] with FillZero → [1,2, 0,0, 3,4]
+        let file = write_list_parquet_with_null_outer_middle();
+        let mut reader =
+            ParquetReader::<f64>::new(file.path(), None, 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        assert_eq!(data, vec![1.0, 2.0, 0.0, 0.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 3);
+        assert_eq!(sample_size, 2);
+    }
+
+    #[test]
+    fn test_parquet_reader_null_outer_row_middle_reject() {
+        // [[1,2], null, [3,4]] with Reject → error
+        let file = write_list_parquet_with_null_outer_middle();
+        let mut reader =
+            ParquetReader::<f64>::new(file.path(), None, 
NullHandling::Reject).unwrap();
+        let result = reader.read_batch();
+        assert!(result.is_err());
+        assert!(result.unwrap_err().to_string().contains("Null outer row"));
+    }
+
+    #[test]
+    fn test_parquet_reader_null_outer_row_first_fill_zero() {
+        // [null, [1,2], [3,4]] — null at row 0 must not corrupt sample_size
+        let file = write_list_parquet_with_null_outer_first();
+        let mut reader =
+            ParquetReader::<f64>::new(file.path(), None, 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        assert_eq!(data, vec![0.0, 0.0, 1.0, 2.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 3);
+        assert_eq!(sample_size, 2);
+    }
+
+    #[test]
+    fn test_parquet_streaming_reader_null_outer_row_middle_fill_zero() {
+        // [[1,2], null, [3,4]] with FillZero → [1,2, 0,0, 3,4]
+        let file = write_list_parquet_with_null_outer_middle();
+        let mut reader =
+            ParquetStreamingReader::<f64>::new(file.path(), None, 
NullHandling::FillZero).unwrap();
+        let mut buffer = vec![0.0_f64; 16];
+        let written = reader.read_chunk(&mut buffer).unwrap();
+        assert_eq!(&buffer[..written], &[1.0, 2.0, 0.0, 0.0, 3.0, 4.0]);
+    }
+
+    #[test]
+    fn test_parquet_streaming_reader_null_outer_row_middle_reject() {
+        // [[1,2], null, [3,4]] with Reject → error
+        let file = write_list_parquet_with_null_outer_middle();
+        let mut reader =
+            ParquetStreamingReader::<f64>::new(file.path(), None, 
NullHandling::Reject).unwrap();
+        let mut buffer = vec![0.0_f64; 16];
+        let result = reader.read_chunk(&mut buffer);
+        assert!(result.is_err());
+        assert!(result.unwrap_err().to_string().contains("Null outer row"));
+    }
+
+    #[test]
+    fn test_parquet_streaming_reader_null_outer_row_first_fill_zero() {
+        // [null, [1,2], [3,4]] — null at row 0 must not seed sample_size to 0
+        let file = write_list_parquet_with_null_outer_first();
+        let mut reader =
+            ParquetStreamingReader::<f64>::new(file.path(), None, 
NullHandling::FillZero).unwrap();
+        let mut buffer = vec![0.0_f64; 16];
+        let written = reader.read_chunk(&mut buffer).unwrap();
+        assert_eq!(&buffer[..written], &[0.0, 0.0, 1.0, 2.0, 3.0, 4.0]);
+    }
+
+    #[test]
+    fn test_parquet_reader_cross_batch_all_null_first_fill_zero() {
+        // [null, null, [1,2], [3,4]] read with batch_size=2:
+        //   batch 1 = [null, null]   — sample_size unknown; must be skipped, 
not counted
+        //   batch 2 = [[1,2], [3,4]] — sample_size established here
+        // Verifies that num_samples is not corrupted by the all-null leading 
batch.
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.append(false); // null row 0
+        builder.append(false); // null row 1
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        let file = write_test_parquet(schema, vec![array]);
+        // batch_size=2 splits into two batches: [null,null] then 
[[1,2],[3,4]].
+        let mut reader =
+            ParquetReader::<f64>::new(file.path(), Some(2), 
NullHandling::FillZero).unwrap();
+        let (data, num_samples, sample_size) = reader.read_batch().unwrap();
+        // All-null first batch is skipped (sample_size unknown → no zeros, no 
count).
+        assert_eq!(data, vec![1.0, 2.0, 3.0, 4.0]);
+        assert_eq!(num_samples, 2);
+        assert_eq!(sample_size, 2);
+    }
+
+    #[test]
+    fn test_parquet_streaming_reader_cross_batch_all_null_first_reject() {
+        // [null, null, [1,2], [3,4]] with batch_size=2, Reject:
+        //   batch 1 = [null, null] — sample_size unknown; must error, not 
silently skip
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.append(false);
+        builder.append(false);
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        let file = write_test_parquet(schema, vec![array]);
+        let mut reader =
+            ParquetStreamingReader::<f64>::new(file.path(), Some(2), 
NullHandling::Reject).unwrap();
+        let mut buffer = vec![0.0_f64; 16];
+        let result = reader.read_chunk(&mut buffer);
+        assert!(result.is_err());
+        assert!(result.unwrap_err().to_string().contains("Null outer row"));
+    }
+
+    #[test]
+    fn test_parquet_streaming_reader_cross_batch_all_null_first_fill_zero() {
+        // [null, null, [1,2], [3,4]] with batch_size=2, FillZero:
+        //   batch 1 = [null, null]   — sample_size unknown; skipped
+        //   batch 2 = [[1,2], [3,4]] — sample_size established; data written
+        let item_field = Arc::new(Field::new("item", DataType::Float64, true));
+        let list_field = Field::new("data", 
DataType::List(item_field.clone()), true);
+        let schema = Arc::new(Schema::new(vec![list_field]));
+
+        let mut builder = ListBuilder::new(Float64Builder::new());
+        builder.append(false);
+        builder.append(false);
+        builder.values().append_slice(&[1.0, 2.0]);
+        builder.append(true);
+        builder.values().append_slice(&[3.0, 4.0]);
+        builder.append(true);
+        let array = Arc::new(builder.finish()) as ArrayRef;
+
+        let file = write_test_parquet(schema, vec![array]);
+        let mut reader =
+            ParquetStreamingReader::<f64>::new(file.path(), Some(2), 
NullHandling::FillZero)
+                .unwrap();
+        let mut buffer = vec![0.0_f64; 16];
+        let written = reader.read_chunk(&mut buffer).unwrap();
+        assert_eq!(&buffer[..written], &[1.0, 2.0, 3.0, 4.0]);
+    }
 }

Reply via email to