From 35f4bd189b299f3d393c01441aee84c176df63dc Mon Sep 17 00:00:00 2001 From: Dhruv Gupta Date: Fri, 2 Oct 2026 14:09:24 -0400 Subject: [PATCH] fix(io): keep null lists of structs when projecting Arrow data The projection visitor rebuilds a list of structs from its offsets alone, which cannot represent a null list, so every null list> was turned into an empty list. On the write path this is persisted into the Parquet file and the null is lost for every reader. Rebuild the list from the original validity and offsets buffers instead, which also keeps the original array offset so sliced batches stay correct. Closes #3833 --- pyiceberg/io/pyarrow.py | 10 ++++++++- tests/integration/test_reads.py | 5 +---- tests/io/test_pyarrow.py | 37 ++++++++++++++++++++++++++++++++- 3 files changed, 46 insertions(+), 6 deletions(-) diff --git a/pyiceberg/io/pyarrow.py b/pyiceberg/io/pyarrow.py index 2dcb8a5795..c9fa0d9a35 100644 --- a/pyiceberg/io/pyarrow.py +++ b/pyiceberg/io/pyarrow.py @@ -2097,7 +2097,15 @@ def list(self, list_type: ListType, list_array: pa.Array | None, value_array: pa if isinstance(value_array, pa.StructArray): # This can be removed once this has been fixed: # https://github.com/apache/arrow/issues/38809 - list_array = pa.LargeListArray.from_arrays(list_array.offsets, value_array) + # Keep the validity and offsets buffers of the original array, otherwise null lists become empty lists + list_array = pa.Array.from_buffers( + list_initializer(value_array.type), + len(list_array), + list_array.buffers()[:2], + list_array.null_count, + list_array.offset, + [value_array], + ) value_array = self._cast_if_needed(list_type.element_field, value_array) arrow_field = list_initializer(self._construct_field(list_type.element_field, value_array.type)) return list_array.cast(arrow_field) diff --git a/tests/integration/test_reads.py b/tests/integration/test_reads.py index a151d62b82..3e07c87b09 100644 --- a/tests/integration/test_reads.py +++ b/tests/integration/test_reads.py @@ -1001,10 +1001,7 @@ def test_null_list_and_map(catalog: Catalog) -> None: arrow_table = table_test_empty_list_and_map.scan().to_arrow() assert arrow_table["col_list"].to_pylist() == [None, []] assert arrow_table["col_map"].to_pylist() == [None, []] - # This should be: - # assert arrow_table["col_list_with_struct"].to_pylist() == [None, [{'test': 1}]] - # Once https://github.com/apache/arrow/issues/38809 has been fixed - assert arrow_table["col_list_with_struct"].to_pylist() == [[], [{"test": 1}]] + assert arrow_table["col_list_with_struct"].to_pylist() == [None, [{"test": 1}]] @pytest.mark.integration diff --git a/tests/io/test_pyarrow.py b/tests/io/test_pyarrow.py index 4d5d4431cb..3f859e0013 100644 --- a/tests/io/test_pyarrow.py +++ b/tests/io/test_pyarrow.py @@ -21,7 +21,7 @@ import tempfile import uuid import warnings -from collections.abc import Iterator +from collections.abc import Callable, Iterator from datetime import date, datetime, timezone from pathlib import Path from typing import Any @@ -3272,6 +3272,41 @@ def test__to_requested_schema_float_promotion( assert result.column(0).to_pylist() == [1.5, 2.25, 3.0, None] +@pytest.mark.parametrize("list_type", [pa.list_, pa.large_list]) +def test__to_requested_schema_null_list_of_structs(list_type: Callable[[pa.DataType], pa.DataType]) -> None: + requested_schema = Schema( + NestedField( + 1, + "col_list_with_struct", + ListType(11, StructType(NestedField(111, "test", IntegerType(), required=False)), element_required=False), + required=False, + ), + NestedField(2, "col_list", ListType(21, IntegerType(), element_required=False), required=False), + ) + arrow_schema = pa.schema( + [ + pa.field("col_list_with_struct", list_type(pa.struct([pa.field("test", pa.int32())]))), + pa.field("col_list", list_type(pa.int32())), + ] + ) + batch = pa.RecordBatch.from_arrays( + [ + pa.array([[{"test": 1}], [], None, [{"test": 2}, None]], type=arrow_schema.field(0).type), + pa.array([[1], [], None, [2, None]], type=arrow_schema.field(1).type), + ], + schema=arrow_schema, + ) + + result = _to_requested_schema(requested_schema, requested_schema, batch) + assert result.column(0).to_pylist() == [[{"test": 1}], [], None, [{"test": 2}, None]] + assert result.column(1).to_pylist() == [[1], [], None, [2, None]] + + # A sliced batch has to keep the nulls as well + result = _to_requested_schema(requested_schema, requested_schema, batch.slice(1, 3)) + assert result.column(0).to_pylist() == [[], None, [{"test": 2}, None]] + assert result.column(1).to_pylist() == [[], None, [2, None]] + + def test_pyarrow_file_io_fs_by_scheme_cache() -> None: # It's better to set up multi-region minio servers for an integration test once `endpoint_url` argument # becomes available for `resolve_s3_region`