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`