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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 1aa0adb717 [python] Read Parquet row windows with OffsetIndex (#9850)
1aa0adb717 is described below

commit 1aa0adb7177511e00a1fe45812e6df8c6bce0021
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Sep 21 11:45:44 2026 +0800

    [python] Read Parquet row windows with OffsetIndex (#9850)
---
 paimon-python/README.md                            |  13 +
 .../pypaimon/common/options/core_options.py        |  16 +
 .../pypaimon/read/reader/format_pyarrow_reader.py  |  66 +-
 .../read/reader/parquet_page_index_reader.py       | 731 +++++++++++++++++++++
 .../pypaimon/tests/parquet_page_index_test.py      | 566 ++++++++++++++++
 5 files changed, 1390 insertions(+), 2 deletions(-)

diff --git a/paimon-python/README.md b/paimon-python/README.md
index 4fa2d7246f..7db4c25386 100644
--- a/paimon-python/README.md
+++ b/paimon-python/README.md
@@ -105,6 +105,19 @@ pip3 install dist/*.tar.gz
 
 The command will install the package and core dependencies to your local 
Python environment.
 
+# Parquet page-index reads
+
+For row-tracking tables with a Parquet OffsetIndex, PyPaimon can read a
+contiguous `_ROW_ID` range without decoding the full row group. This is enabled
+by default and can be disabled with the table option:
+
+```python
+table = table.copy({"parquet.filter.columnindex.enabled": "false"})
+```
+
+Unsupported reads use the normal path. Reading fewer bytes may require more
+object-store requests.
+
 # Native scan planning
 
 PyPaimon can plan splits with the optional `pypaimon-rust` package while 
retaining
diff --git a/paimon-python/pypaimon/common/options/core_options.py 
b/paimon-python/pypaimon/common/options/core_options.py
index 5d4cfa34dd..79b0202125 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -1158,6 +1158,19 @@ class CoreOptions:
         .with_description("Read batch size for any file format if it 
supports.")
     )
 
+    PARQUET_COLUMN_INDEX_ENABLED: ConfigOption[bool] = (
+        ConfigOptions.key("parquet.filter.columnindex.enabled")
+        .boolean_type()
+        .default_value(True)
+        .with_description(
+            "Enable Parquet page-index pruning. PyPaimon currently uses 
OffsetIndex "
+            "metadata for contiguous row windows. "
+            "Requires existing offset indexes; nested fields use common leaf 
row boundaries. "
+            "Unsupported or expensive selections use the ordinary reader. "
+            "Does not enable ColumnIndex predicate filtering."
+        )
+    )
+
     READ_PARALLELISM: ConfigOption[int] = (
         ConfigOptions.key("read.parallelism")
         .int_type()
@@ -1917,6 +1930,9 @@ class CoreOptions:
     def read_batch_size(self, default=None) -> int:
         return self.options.get(CoreOptions.READ_BATCH_SIZE, default or 1024)
 
+    def parquet_column_index_enabled(self) -> bool:
+        return self.options.get(CoreOptions.PARQUET_COLUMN_INDEX_ENABLED)
+
     def read_parallelism(self, default=None) -> Optional[int]:
         return self.options.get(CoreOptions.READ_PARALLELISM, default)
 
diff --git a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py 
b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
index ead8dd3bdf..df9e3e7002 100644
--- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
+++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
@@ -525,6 +525,8 @@ class FormatPyArrowReader(RecordBatchReader):
 
         # Read projected VARIANT columns in bounded batches.
         self._parquet_file = None
+        self._parquet_source = None
+        self._page_index_reader = None
         self._orc_file = None
         self._orc_source = None
         if (self._bounded_variant_read
@@ -533,8 +535,25 @@ class FormatPyArrowReader(RecordBatchReader):
                     and self._selected_shared_map_paths)):
             import pyarrow.parquet as pq
             # ParquetFile(filesystem=...) is unavailable in PyArrow 6.
-            self._parquet_file = pq.ParquetFile(
-                file_io.filesystem.open_input_file(file_path_for_pyarrow))
+            self._parquet_source = file_io.filesystem.open_input_file(
+                file_path_for_pyarrow)
+            try:
+                self._parquet_file = pq.ParquetFile(self._parquet_source)
+                if (self._selected_parquet_row_groups is not None
+                        and options is not None
+                        and options.parquet_column_index_enabled()
+                        and self._row_group_cache is None
+                        and not self._bounded_variant_read):
+                    from pypaimon.read.reader.parquet_page_index_reader import 
(
+                        ParquetPageIndexReader,
+                    )
+                    self._page_index_reader = ParquetPageIndexReader.create(
+                        self._parquet_source, self._parquet_file,
+                        self._row_group_read_columns(),
+                        self._selected_parquet_row_groups, batch_size)
+            except BaseException:
+                self._parquet_source.close()
+                raise
         if file_format == 'orc' and self._selected_shared_map_paths:
             import pyarrow.orc as orc
             self._orc_source = file_io.filesystem.open_input_file(
@@ -542,6 +561,11 @@ class FormatPyArrowReader(RecordBatchReader):
             self._orc_file = orc.ORCFile(self._orc_source)
         if self._exhausted:
             self._raw_batches = iter(())
+        elif self._page_index_reader is not None:
+            # Page selection already preserves original row positions. Slice
+            # fallback row groups here too, before mixing the two streams.
+            self._range_slicer = None
+            self._raw_batches = self._iter_page_index_batches(selected_infos, 
runs)
         elif self._parquet_file is not None:
             self._raw_batches = self._iter_row_group_batches()
         elif self._orc_file is not None:
@@ -618,6 +642,37 @@ class FormatPyArrowReader(RecordBatchReader):
                     if out.num_rows:
                         yield out
 
+    def _iter_page_index_batches(self, selected_infos, runs):
+        select = self._select_nested_fields if self._has_nested_path else 
self._select_existing_fields
+        run_index = 0
+        for group, (offset, count) in zip(
+                self._selected_parquet_row_groups, selected_infos):
+            while run_index < len(runs) and runs[run_index][1] < offset:
+                run_index += 1
+            local_runs = []
+            position = run_index
+            while position < len(runs) and runs[position][0] < offset + count:
+                lower, upper = runs[position]
+                local_runs.append((max(0, lower - offset),
+                                   min(count - 1, upper - offset)))
+                position += 1
+            batches = self._page_index_reader.read_row_group(group, local_runs)
+            if batches is None:
+                raw = self._read_parquet_row_group_batches(
+                    group, self._row_group_read_columns())
+                slicer = _RowRunSlicer([(0, count)], local_runs)
+                while True:
+                    batch = slicer.next_batch(raw)
+                    if batch is None:
+                        break
+                    yield select(batch)
+            else:
+                try:
+                    for batch in batches:
+                        yield select(batch)
+                finally:
+                    batches.close()
+
     def _read_parquet_row_group_batches(self, row_group, columns):
         return self._parquet_file.iter_batches(
             row_groups=[row_group],
@@ -859,12 +914,19 @@ class FormatPyArrowReader(RecordBatchReader):
         return batch
 
     def close(self):
+        close_batches = getattr(self._raw_batches, 'close', None)
+        if close_batches is not None:
+            close_batches()
         self._raw_batches = None
         if self._parquet_file is not None:
             close = getattr(self._parquet_file, 'close', None)
             if close is not None:
                 close()
             self._parquet_file = None
+        if self._parquet_source is not None:
+            self._parquet_source.close()
+            self._parquet_source = None
+            self._page_index_reader = None
         if self._orc_source is not None:
             self._orc_source.close()
             self._orc_source = None
diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py 
b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py
new file mode 100644
index 0000000000..6ca5ca25ca
--- /dev/null
+++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py
@@ -0,0 +1,731 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+"""Read contiguous Parquet row windows using their OffsetIndex.
+
+Selected encoded pages are placed in bounded, in-memory Parquet files. PyArrow
+still decodes the pages, including dictionary and compression encodings. Source
+files are never rewritten. Nested fields retain their complete physical schema
+and align their leaf columns at common row boundaries. Files without indexes 
and
+disjoint ranges use the normal reader: per-page seeks can amplify requests in
+filesystems that prefetch remote data (including Jindo).
+"""
+
+import base64
+import bisect
+import struct
+
+import pyarrow as pa
+import pyarrow.parquet as pq
+
+
+# Bound encoded page data retained by the temporary column files.
+_MAX_PAGE_BYTES = 32 * 1024 * 1024
+# Bound the generic FileMetaData tree before creating Python objects for it.
+_MAX_FOOTER_BYTES = 1024 * 1024
+_MAX_FOOTER_COLUMN_CHUNKS = 1024
+_MAX_FOOTER_ITEMS = 64 * 1024
+# Bound all serialized OffsetIndexes and their retained typed PageLocations.
+_MAX_INDEX_BYTES = 8 * 1024 * 1024
+_MAX_PAGE_LOCATIONS = 128 * 1024
+_MAX_PAGE_HEADER_ITEMS = 4096
+
+
+class _PageIndexBudgetExceeded(Exception):
+    pass
+
+
+class _Compact:
+    """Thrift compact values used by Parquet metadata (no generated 
bindings)."""
+
+    def __init__(self, data, max_items=None):
+        self.data = memoryview(data)
+        self.position = 0
+        self.remaining_items = max_items
+
+    def consume(self, count):
+        if self.remaining_items is not None:
+            if count > self.remaining_items:
+                raise _PageIndexBudgetExceeded(
+                    "Parquet compact metadata exceeds object budget")
+            self.remaining_items -= count
+
+    def take(self, size):
+        end = self.position + size
+        if size < 0 or end > len(self.data):
+            raise ValueError("Truncated Parquet page-index metadata")
+        result = self.data[self.position:end]
+        self.position = end
+        return result
+
+    def unsigned(self):
+        result = 0
+        for shift in range(0, 70, 7):
+            value = self.take(1)[0]
+            result |= (value & 127) << shift
+            if value < 128:
+                return result
+        raise ValueError("Invalid Parquet compact integer")
+
+    def value(self, kind, depth=0):
+        if depth > 64:
+            raise ValueError("Parquet metadata nesting exceeds 64 levels")
+        if kind in (1, 2):
+            return kind == 1
+        if kind == 3:
+            return self.take(1).tobytes()
+        if kind in (4, 5, 6):
+            value = self.unsigned()
+            return (value >> 1) ^ -(value & 1)
+        if kind == 7:
+            return self.take(8).tobytes()
+        if kind == 8:
+            return self.take(self.unsigned()).tobytes()
+        if kind in (9, 10):
+            header = self.take(1)[0]
+            count, element = header >> 4, header & 15
+            if count == 15:
+                count = self.unsigned()
+            if count > len(self.data) - self.position:
+                raise ValueError("Invalid Parquet compact collection size")
+            self.consume(count)
+            return element, [
+                self.value(self.take(1)[0] if element in (1, 2) else element, 
depth + 1)
+                for _ in range(count)
+            ]
+        if kind == 12:
+            fields = {}
+            previous = 0
+            while True:
+                header = self.take(1)[0]
+                if header == 0:
+                    return fields
+                delta, field_kind = header >> 4, header & 15
+                field = previous + delta if delta else self.value(4)
+                if field in fields:
+                    raise ValueError("Duplicate Parquet compact field")
+                self.consume(1)
+                fields[field] = field_kind, self.value(field_kind, depth + 1)
+                previous = field
+        raise ValueError("Unsupported Parquet compact type: {}".format(kind))
+
+
+class _BoundedCompactDecoder:
+    """Skip unknown compact fields without materializing object trees."""
+
+    def __init__(self, data, max_items):
+        self.parser = _Compact(data)
+        self.remaining_items = max_items
+
+    def _field(self, previous):
+        header = self.parser.take(1)[0]
+        if header == 0:
+            return None
+        delta, kind = header >> 4, header & 15
+        field = previous + delta if delta else self.parser.value(4)
+        if field <= 0:
+            raise ValueError("Invalid Parquet compact field")
+        self._consume(1)
+        return field, kind
+
+    def _collection(self):
+        header = self.parser.take(1)[0]
+        count, element = header >> 4, header & 15
+        if count == 15:
+            count = self.parser.unsigned()
+        if count > len(self.parser.data) - self.parser.position:
+            raise ValueError("Invalid Parquet compact collection size")
+        return count, element
+
+    def _consume(self, count):
+        if count > self.remaining_items:
+            raise _PageIndexBudgetExceeded(
+                "Parquet compact metadata exceeds object budget")
+        self.remaining_items -= count
+
+    def _skip_collection_value(self, kind, depth):
+        if kind in (1, 2):
+            actual = self.parser.take(1)[0]
+            if actual not in (1, 2):
+                raise ValueError("Invalid Parquet compact boolean")
+        else:
+            self._skip(kind, depth)
+
+    def _skip(self, kind, depth=0):
+        if depth > 64:
+            raise ValueError("Parquet metadata nesting exceeds 64 levels")
+        if kind in (1, 2):
+            return
+        if kind == 3:
+            self.parser.take(1)
+            return
+        if kind in (4, 5, 6):
+            self.parser.unsigned()
+            return
+        if kind == 7:
+            self.parser.take(8)
+            return
+        if kind == 8:
+            self.parser.take(self.parser.unsigned())
+            return
+        if kind in (9, 10):
+            count, element = self._collection()
+            self._consume(count)
+            for _ in range(count):
+                self._skip_collection_value(element, depth + 1)
+            return
+        if kind == 11:
+            count = self.parser.unsigned()
+            self._consume(count * 2)
+            if count:
+                kinds = self.parser.take(1)[0]
+                key_kind, value_kind = kinds >> 4, kinds & 15
+                for _ in range(count):
+                    self._skip_collection_value(key_kind, depth + 1)
+                    self._skip_collection_value(value_kind, depth + 1)
+            return
+        if kind == 12:
+            previous = 0
+            while True:
+                field = self._field(previous)
+                if field is None:
+                    return
+                field_id, field_kind = field
+                previous = field_id
+                self._skip(field_kind, depth + 1)
+        raise ValueError("Unsupported Parquet compact type: {}".format(kind))
+
+
+class _OffsetIndexDecoder(_BoundedCompactDecoder):
+    """Decode typed PageLocations without materializing generic Thrift 
trees."""
+
+    def __init__(self, data, max_locations):
+        # PageLocations retain three fields each. Allow the standard optional
+        # per-page byte-count list plus a small amount of forward metadata.
+        super().__init__(data, max_locations * 6 + 16)
+        self.max_locations = max_locations
+
+    def decode(self):
+        locations = None
+        previous = 0
+        while True:
+            field = self._field(previous)
+            if field is None:
+                break
+            field_id, kind = field
+            previous = field_id
+            if field_id == 1:
+                if locations is not None:
+                    raise ValueError("Duplicate Parquet OffsetIndex field")
+                if kind != 9:
+                    raise ValueError("Invalid Parquet OffsetIndex page 
locations")
+                locations = self._locations()
+            else:
+                self._skip(kind)
+        if locations is None or self.parser.position != len(self.parser.data):
+            raise ValueError("Invalid Parquet OffsetIndex")
+        return locations
+
+    def _locations(self):
+        count, element = self._collection()
+        if element != 12:
+            raise ValueError("Invalid Parquet OffsetIndex page locations")
+        if count > self.max_locations:
+            raise _PageIndexBudgetExceeded(
+                "Parquet OffsetIndex exceeds page-location budget")
+        self._consume(count)
+        return [self._location() for _ in range(count)]
+
+    def _location(self):
+        values = [None, None, None]
+        expected = (6, 5, 6)
+        previous = 0
+        while True:
+            field = self._field(previous)
+            if field is None:
+                break
+            field_id, kind = field
+            previous = field_id
+            if 1 <= field_id <= 3:
+                if values[field_id - 1] is not None:
+                    raise ValueError("Duplicate Parquet PageLocation field")
+                if kind != expected[field_id - 1]:
+                    raise ValueError("Invalid Parquet PageLocation field type")
+                values[field_id - 1] = self.parser.value(kind)
+            else:
+                self._skip(kind)
+        if any(value is None for value in values):
+            raise ValueError("Missing Parquet PageLocation field")
+        offset, size, first_row = values
+        if offset < 0 or size <= 0 or first_row < 0:
+            raise ValueError("Invalid Parquet PageLocation")
+        return offset, size, first_row
+
+
+class _PageHeaderDecoder(_BoundedCompactDecoder):
+    """Decode only PageHeader fields required to validate selected pages."""
+
+    _NESTED_FIELDS = {5: (1,), 7: (1,), 8: (1, 3)}
+
+    def __init__(self, data):
+        super().__init__(data, _MAX_PAGE_HEADER_ITEMS)
+
+    def decode(self):
+        result = {}
+        previous = 0
+        while True:
+            field = self._field(previous)
+            if field is None:
+                break
+            field_id, kind = field
+            previous = field_id
+            if field_id in (1, 2, 3):
+                if field_id in result:
+                    raise ValueError("Duplicate Parquet PageHeader field")
+                if kind != 5:
+                    raise ValueError("Invalid Parquet PageHeader field type")
+                result[field_id] = kind, self.parser.value(kind)
+            elif field_id in self._NESTED_FIELDS:
+                if field_id in result:
+                    raise ValueError("Duplicate Parquet PageHeader field")
+                if kind != 12:
+                    raise ValueError("Invalid Parquet PageHeader field type")
+                result[field_id] = kind, self._integer_struct(
+                    self._NESTED_FIELDS[field_id])
+            else:
+                self._skip(kind)
+        if any(field not in result for field in (1, 2, 3)):
+            raise ValueError("Missing Parquet PageHeader field")
+        return result, self.parser.position
+
+    def _integer_struct(self, required):
+        result = {}
+        previous = 0
+        while True:
+            field = self._field(previous)
+            if field is None:
+                break
+            field_id, kind = field
+            previous = field_id
+            if field_id in required:
+                if field_id in result:
+                    raise ValueError("Duplicate Parquet page header field")
+                if kind != 5:
+                    raise ValueError("Invalid Parquet page header field type")
+                result[field_id] = kind, self.parser.value(kind)
+            else:
+                self._skip(kind)
+        if any(field not in result for field in required):
+            raise ValueError("Missing Parquet page header field")
+        return result
+
+
+def _decode_offset_index(data, max_locations):
+    return _OffsetIndexDecoder(data, max_locations).decode()
+
+
+def _decode_page_header(data):
+    return _PageHeaderDecoder(data).decode()
+
+
+def _unsigned(value):
+    result = bytearray()
+    while value >= 128:
+        result.append((value & 127) | 128)
+        value >>= 7
+    result.append(value)
+    return bytes(result)
+
+
+def _encode(kind, value):
+    if kind in (1, 2):
+        return bytes([1 if value else 2])
+    if kind in (3, 7):
+        return value
+    if kind in (4, 5, 6):
+        return _unsigned(value * 2 if value >= 0 else -value * 2 - 1)
+    if kind == 8:
+        return _unsigned(len(value)) + value
+    if kind in (9, 10):
+        element, items = value
+        size = len(items)
+        header = bytes([(min(size, 15) << 4) | element])
+        if size >= 15:
+            header += _unsigned(size)
+        return header + b"".join(_encode(element, item) for item in items)
+    if kind == 12:
+        result = bytearray()
+        previous = 0
+        for field, (field_kind, item) in sorted(value.items()):
+            delta = field - previous
+            if field_kind in (1, 2):
+                field_kind = 1 if item else 2
+            if 0 < delta < 16:
+                result.append((delta << 4) | field_kind)
+            else:
+                result.append(field_kind)
+                result.extend(_encode(4, field))
+            if field_kind not in (1, 2):
+                result.extend(_encode(field_kind, item))
+            previous = field
+        result.append(0)
+        return bytes(result)
+    raise ValueError("Unsupported Parquet compact type: {}".format(kind))
+
+
+def _get(fields, field, default=None):
+    return fields[field][1] if field in fields else default
+
+
+def _read_exact(source, offset, length):
+    if offset < 4 or length <= 0:
+        raise ValueError("Invalid Parquet page-index byte range")
+    data = source.read_at(length, offset)
+    if len(data) != length:
+        raise OSError("Truncated Parquet page-index byte range")
+    return data
+
+
+def _read_index_ranges(source, ranges):
+    groups = []
+    for key, offset, length in sorted(ranges, key=lambda item: item[1]):
+        if offset < 4 or length <= 0:
+            raise ValueError("Invalid Parquet page-index byte range")
+        end = offset + length
+        if groups and offset < groups[-1][1]:
+            raise ValueError("Overlapping Parquet page-index byte ranges")
+        if groups and offset == groups[-1][1]:
+            groups[-1][1] = end
+            groups[-1][2].append((key, offset, length))
+        else:
+            groups.append([offset, end, [(key, offset, length)]])
+    result = {}
+    for start, end, members in groups:
+        data = memoryview(_read_exact(source, start, end - start))
+        for key, offset, length in members:
+            result[key] = data[offset - start:offset - start + length]
+    return result
+
+
+class ParquetPageIndexReader:
+    def __init__(self, source, metadata, schema, footer, columns, fields, 
batch_size):
+        self.source = source
+        self.metadata = metadata
+        self.schema = schema
+        self.footer = footer
+        self.columns = columns
+        self.fields = fields
+        self.batch_size = batch_size
+
+    @classmethod
+    def create(cls, source, parquet_file, columns, row_groups, batch_size):
+        metadata = parquet_file.metadata
+        schema = parquet_file.schema_arrow
+        if not columns or len(set(schema.names)) != len(schema):
+            return None
+        indices = [schema.get_field_index(name) for name in columns]
+        if any(index < 0 for index in indices) or len(set(indices)) != 
len(indices):
+            return None
+        if not any(getattr(metadata.row_group(group).column(index),
+                           "has_offset_index", False)
+                   for group in row_groups for index in 
range(metadata.num_columns)):
+            return None
+        if (metadata.serialized_size > _MAX_FOOTER_BYTES
+                or metadata.num_row_groups * metadata.num_columns
+                > _MAX_FOOTER_COLUMN_CHUNKS):
+            return None
+        output = pa.BufferOutputStream()
+        metadata.write_metadata_file(output)
+        serialized = output.getvalue().to_pybytes()
+        length = struct.unpack("<I", serialized[-8:-4])[0]
+        if length > _MAX_FOOTER_BYTES:
+            return None
+        try:
+            footer = _Compact(
+                serialized[-8 - length:-8], _MAX_FOOTER_ITEMS).value(12)
+        except _PageIndexBudgetExceeded:
+            return None
+        if 8 in footer or 9 in footer:
+            return None  # Encrypted pages need the original file identity/AAD.
+        elements = _get(footer, 2)[1]
+        # Parquet stores a preorder schema tree and one chunk per physical 
leaf.
+        # Arrow field positions cannot be used as physical column positions.
+        fields = []
+        position, leaf = 1, 0
+        for _ in range(_get(elements[0], 5)):
+            start, first_leaf, pending = position, leaf, 1
+            while pending:
+                if position >= len(elements):
+                    raise ValueError("Truncated Parquet schema tree")
+                element = elements[position]
+                children = _get(element, 5, 0)
+                if children < 0 or (1 in element and children) or (1 not in 
element and not children):
+                    raise ValueError("Invalid Parquet schema child count")
+                pending += children - 1
+                leaf += int(1 in element)
+                position += 1
+            fields.append((elements[start:position], list(range(first_leaf, 
leaf))))
+        if (position != len(elements) or leaf != metadata.num_columns
+                or len(fields) != len(schema)):
+            return None
+        if not any(all(getattr(metadata.row_group(group).column(leaf),
+                               "has_offset_index", False)
+                       for index in indices for leaf in fields[index][1])
+                   for group in row_groups):
+            return None
+        return cls(source, metadata, schema, footer, indices, fields, 
batch_size)
+
+    def read_row_group(self, group, runs):
+        """Return selected batches, or None when the ordinary reader is 
cheaper."""
+        # Decide before reading indexes so scattered selections preserve the
+        # existing I/O pattern. A future multi-range path needs an I/O planner
+        # that accounts for filesystem prefetch, not just compressed page 
sizes.
+        if len(runs) != 1:
+            return None
+        row_group = _get(self.footer, 4)[1][group]
+        row_count = _get(row_group, 3)
+        if sum(upper - lower + 1 for lower, upper in runs) >= row_count:
+            return None
+        chunks = _get(row_group, 1)[1]
+        physical_columns = [leaf for index in self.columns for leaf in 
self.fields[index][1]]
+        if any(4 not in chunks[index] or 5 not in chunks[index]
+               or _get(chunks[index], 1) or 8 in chunks[index] or 9 in 
chunks[index]
+               or 10 in _get(chunks[index], 3)  # Legacy index pages.
+               for index in physical_columns):
+            return None
+        indexed = {}
+        plans = []
+        selected_bytes = 0
+        full_bytes = 0
+        index_ranges = []
+        for index in sorted(physical_columns):
+            chunk = chunks[index]
+            index_size = _get(chunk, 5)
+            index_ranges.append((index, _get(chunk, 4), index_size))
+        index_bytes = sum(length for _, _, length in index_ranges)
+        if index_bytes > _MAX_INDEX_BYTES:
+            return None
+        raw_indexes = _read_index_ranges(self.source, index_ranges)
+        remaining_locations = _MAX_PAGE_LOCATIONS
+        for index in sorted(physical_columns):
+            chunk = chunks[index]
+            try:
+                locations = _decode_offset_index(raw_indexes[index], 
remaining_locations)
+            except _PageIndexBudgetExceeded:
+                return None
+            remaining_locations -= len(locations)
+            column = _get(chunk, 3)
+            data_offset = _get(column, 9)
+            dictionary_offset = _get(column, 11, data_offset)
+            chunk_end = dictionary_offset + _get(column, 7)
+            starts = [page[2] for page in locations]
+            if not starts or starts[0] != 0 or starts[-1] >= row_count:
+                raise ValueError("Invalid Parquet OffsetIndex row boundaries")
+            previous_end = data_offset
+            previous_row = -1
+            for page in locations:
+                offset, size, first_row = page
+                if (offset < previous_end or size <= 0 or offset + size > 
chunk_end
+                        or first_row <= previous_row):
+                    raise ValueError("Invalid Parquet OffsetIndex page 
location")
+                previous_end, previous_row = offset + size, first_row
+            if locations[0][0] != data_offset or dictionary_offset > 
data_offset:
+                raise ValueError("Invalid Parquet OffsetIndex first page")
+            indexed[index] = (column, dictionary_offset, data_offset - 
dictionary_offset,
+                              locations, starts)
+            full_bytes += _get(column, 7)
+        for field in sorted(self.columns):
+            leaves = self.fields[field][1]
+            # OffsetIndex pages must start at row boundaries (repetition level 
0).
+            # Keep all leaves of a field aligned so Arrow can reconstruct 
nesting.
+            # ponytail: common boundaries may widen to the whole group; 
independent
+            # leaf decoding/reassembly can recover savings if this becomes 
costly.
+            boundaries = set(indexed[leaves[0]][4])
+            for leaf in leaves[1:]:
+                boundaries.intersection_update(indexed[leaf][4])
+            boundaries = sorted(boundaries) + [row_count]
+            lower, upper = runs[0]
+            lower = boundaries[bisect.bisect_right(boundaries, lower) - 1]
+            end = boundaries[bisect.bisect_right(boundaries, upper)]
+            column_plans = []
+            for index in leaves:
+                column, dictionary_offset, dictionary_size, locations, starts 
= indexed[index]
+                selected = range(bisect.bisect_left(starts, lower),
+                                 bisect.bisect_left(starts, end))
+                pages, infos = [], []
+                for position in selected:
+                    page = locations[position]
+                    pages.append(page[:2])
+                    next_row = starts[position + 1] if position + 1 < 
len(starts) else row_count
+                    infos.append((starts[position], next_row - 
starts[position]))
+                selected_bytes += dictionary_size + sum(size for _, size in 
pages)
+                column_plans.append((index, column, dictionary_offset, 
dictionary_size, pages, infos))
+            plans.append((field, column_plans, [(lower, end - lower)]))
+        if (selected_bytes + index_bytes >= full_bytes
+                or selected_bytes > _MAX_PAGE_BYTES):
+            return None
+        batches = self._batches(plans, runs)
+        try:
+            first = next(batches)
+        except _PageIndexBudgetExceeded:
+            batches.close()
+            return None
+
+        def prepared_batches():
+            try:
+                yield first
+                yield from batches
+            finally:
+                batches.close()
+
+        # Every selected PageHeader is decoded while preparing the first batch,
+        # so a budget fallback cannot duplicate rows already returned to 
callers.
+        return prepared_batches()
+
+    def _column_payload(self, plan):
+        index, column, dictionary_offset, dictionary_size, pages, infos = plan
+        ranges = ([(dictionary_offset, dictionary_size)] if dictionary_size 
else []) + pages
+        # Coalesce adjacent dictionary/data pages without fetching skipped 
pages.
+        groups = []
+        for offset, length in ranges:
+            if groups and groups[-1][0] + groups[-1][1] == offset:
+                groups[-1][1] += length
+            else:
+                groups.append([offset, length])
+        payload = b"".join(_read_exact(self.source, offset, length)
+                           for offset, length in groups)
+        cursor = 0
+        uncompressed_size = 0
+        num_values = 0
+        repeated = self.metadata.schema.column(index).max_repetition_level > 0
+        for position, (_, length) in enumerate(ranges):
+            header, header_size = _decode_page_header(
+                memoryview(payload)[cursor:cursor + length])
+            if header_size + _get(header, 3) != length:
+                raise ValueError("Parquet page size disagrees with 
OffsetIndex")
+            if dictionary_size and position == 0:
+                if _get(header, 1) != 2 or 7 not in header:
+                    raise ValueError("Invalid Parquet dictionary page")
+            else:
+                expected = infos[position - bool(dictionary_size)][1]
+                page_type = _get(header, 1)
+                if page_type == 0:
+                    values = _get(_get(header, 5), 1)
+                    actual = expected if repeated else values
+                elif page_type == 3:
+                    page_header = _get(header, 8)
+                    actual = _get(page_header, 3)
+                    values = _get(page_header, 1)
+                    if not repeated and values != actual:
+                        raise ValueError("Invalid non-repeated Parquet data 
page")
+                else:
+                    raise ValueError("Invalid Parquet data page type")
+                if actual != expected or values < expected:
+                    raise ValueError("Parquet page rows disagree with 
OffsetIndex")
+                num_values += values
+            uncompressed_size += header_size + _get(header, 2)
+            cursor += length
+
+        patched_column = {key: value for key, value in column.items() if key 
<= 8}
+        patched_column.update({5: (6, num_values), 6: (6, uncompressed_size),
+                               7: (6, len(payload)), 9: (6, 4 + 
dictionary_size)})
+        if dictionary_size:
+            patched_column[11] = (6, 4)
+        return payload, patched_column, uncompressed_size
+
+    def _column_batches(self, plan, runs):
+        from pypaimon.read.reader.format_pyarrow_reader import _RowRunSlicer
+
+        index, column_plans, infos = plan
+        payloads, chunks = [], []
+        offset, uncompressed_size = 0, 0
+        for column_plan in column_plans:
+            payload, column, size = self._column_payload(column_plan)
+            for field in (9, 11):
+                if field in column:
+                    column[field] = (6, _get(column, field) + offset)
+            chunks.append({2: (6, 0), 3: (12, column)})
+            payloads.append(payload)
+            offset += len(payload)
+            uncompressed_size += size
+        num_rows = sum(count for _, count in infos)
+        patched_group = {1: (9, (12, chunks)),
+                         2: (6, uncompressed_size), 3: (6, num_rows)}
+        elements = _get(self.footer, 2)[1]
+        root = dict(elements[0])
+        root[5] = (5, 1)
+        schema = pa.schema([self.schema.field(index)])
+        arrow_schema = base64.b64encode(schema.serialize().to_pybytes())
+        footer = {1: self.footer[1], 2: (9, (12, [root] + 
self.fields[index][0])),
+                  3: (6, num_rows), 4: (9, (12, [patched_group])),
+                  5: (9, (12, [{1: (8, b"ARROW:schema"), 2: (8, 
arrow_schema)}]))}
+        if 6 in self.footer:
+            footer[6] = self.footer[6]
+        encoded = _encode(12, footer)
+        data = b"".join([b"PAR1"] + payloads + [encoded, struct.pack("<I", 
len(encoded)), b"PAR1"])
+        del payloads, payload
+        reader = pq.ParquetFile(pa.BufferReader(data))
+        try:
+            def checked_batches():
+                count = 0
+                for batch in reader.iter_batches(batch_size=self.batch_size, 
use_threads=False):
+                    count += batch.num_rows
+                    if count > num_rows:
+                        raise ValueError("Parquet decoded rows disagree with 
OffsetIndex")
+                    yield batch
+                if count != num_rows:
+                    raise ValueError("Parquet decoded rows disagree with 
OffsetIndex")
+
+            batches = checked_batches()
+            slicer = _RowRunSlicer(infos, runs)
+            while True:
+                batch = slicer.next_batch(batches)
+                if batch is None:
+                    break
+                yield batch.column(0)
+        finally:
+            reader.close()
+
+    def _batches(self, plans, runs):
+        readers = [self._column_batches(plan, runs) for plan in plans]
+        remaining = sum(upper - lower + 1 for lower, upper in runs)
+        positions = {plan[0]: position for position, plan in enumerate(plans)}
+        projection = [positions[index] for index in self.columns]
+        try:
+            arrays = [next(reader, None) for reader in readers]
+            offsets = [0] * len(readers)
+            schema = pa.schema([self.schema.field(index) for index in 
self.columns])
+            while any(array is not None for array in arrays):
+                if any(array is None for array in arrays):
+                    raise ValueError("Parquet page-index columns have 
different row counts")
+                count = min(len(array) - offset for array, offset in 
zip(arrays, offsets))
+                remaining -= count
+                if count <= 0 or remaining < 0:
+                    raise ValueError("Invalid Parquet page-index result 
length")
+                yield pa.RecordBatch.from_arrays(
+                    [arrays[index].slice(offsets[index], count) for index in 
projection],
+                    schema=schema)
+                for index, array in enumerate(arrays):
+                    offsets[index] += count
+                    if offsets[index] == len(array):
+                        arrays[index] = next(readers[index], None)
+                        offsets[index] = 0
+            if remaining:
+                raise ValueError("Truncated Parquet page-index result")
+        finally:
+            for reader in readers:
+                reader.close()
diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py 
b/paimon-python/pypaimon/tests/parquet_page_index_test.py
new file mode 100644
index 0000000000..8a33e0b1dc
--- /dev/null
+++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py
@@ -0,0 +1,566 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+import hashlib
+from concurrent.futures import ThreadPoolExecutor
+from decimal import Decimal
+from unittest.mock import patch
+
+import pyarrow as pa
+import pyarrow.fs as pafs
+import pyarrow.parquet as pq
+import pytest
+
+from pypaimon.common.options import Options
+from pypaimon.common.options.core_options import CoreOptions
+from pypaimon.filesystem.local_file_io import LocalFileIO
+from pypaimon.read.reader import format_pyarrow_reader as reader_module
+from pypaimon.read.reader import parquet_page_index_reader as page_module
+from pypaimon.schema.data_types import AtomicType, DataField, 
PyarrowFieldParser
+from pypaimon.tests.parquet_metadata_cache_test import _CountingLocalFileSystem
+
+
+pytestmark = pytest.mark.skipif(
+    int(pa.__version__.split('.')[0]) < 13,
+    reason='Writing page indexes requires PyArrow 13 or later')
+N = 16384
+RUNS = [(0, 2), (125, 132), (4500, 4540), (N - 2, N - 1)]
+FIELDS = [DataField(0, 'id', AtomicType('BIGINT')),
+          DataField(1, 'payload', AtomicType('STRING'))]
+PAGE_INDEX_OPTIONS = 
CoreOptions(Options({'parquet.filter.columnindex.enabled': 'true'}))
+
+
[email protected]
+def fixture(tmp_path):
+    path = str(tmp_path / 'indexed.parquet')
+    table = pa.table({'id': range(N), 'payload': [
+        hashlib.sha256(str(i).encode()).hexdigest() for i in range(N)]})
+    pq.write_table(table, path, write_page_index=True,
+                   data_page_size=4096, write_batch_size=128,
+                   dictionary_pagesize_limit=16 * 1024,
+                   row_group_size=N // 2)
+    counter = _CountingLocalFileSystem(skip_instance_cache=True)
+    file_io = LocalFileIO(str(tmp_path), Options({}))
+    file_io.filesystem = pafs.PyFileSystem(pafs.FSSpecHandler(counter))
+    reader_module._reset_file_format_dataset_cache()
+    yield path, table, file_io, counter
+    reader_module._reset_file_format_dataset_cache()
+
+
+def _expected(table, runs):
+    return table.take(pa.array(sorted({
+        i for lower, upper in runs for i in range(max(0, lower), 
min(len(table), upper + 1))
+    }), type=pa.int64()))
+
+
+def _read(fixture, baseline=False, fields=FIELDS, **kwargs):
+    path, table, file_io, counter = fixture
+    counter.reset_counts()
+    kwargs.setdefault('row_ranges', RUNS)
+    kwargs.setdefault('options', None if baseline else PAGE_INDEX_OPTIONS)
+    reader = reader_module.FormatPyArrowReader(
+        file_io, 'parquet', path, fields, None, batch_size=71, **kwargs)
+    try:
+        batches = []
+        while True:
+            batch = reader.read_arrow_batch()
+            if batch is None:
+                break
+            assert batch.num_rows <= 71
+            batches.append(batch)
+        result = pa.Table.from_batches(batches) if batches else None
+        return result, list(counter.reads)
+    finally:
+        reader.close()
+
+
[email protected]('cache_size', ['0 b', '50 mb'])
+def test_sparse_reads_skip_bytes_and_preserve_results(fixture, cache_size):
+    path, table, file_io, counter = fixture
+    file_io.properties = Options({'file-format.metadata-cache.max-size': 
cache_size})
+    runs = [(4500, 4540)]
+    baseline, baseline_reads = _read(fixture, baseline=True, row_ranges=runs)
+    reader_module._reset_file_format_dataset_cache()
+    for _ in range(2):  # Cold and warm metadata cache.
+        result, reads = _read(fixture, row_ranges=runs)
+        assert result.equals(baseline)
+        assert result.equals(_expected(table, runs))
+        assert sum(size for _, size in reads) < sum(size for _, size in 
baseline_reads) / 2
+
+
+def test_ranges_projection_missing_fields_and_fallback_in_same_file(fixture):
+    # First row group is read in full; the second uses selected pages.
+    runs = [(-10, N // 2 - 1), (N - 2, N + 10)]
+    fields = [FIELDS[1], DataField(2, 'added', AtomicType('INT')), FIELDS[0]]
+    baseline, _ = _read(fixture, baseline=True, fields=fields, row_ranges=runs)
+    result, _ = _read(fixture, fields=fields, row_ranges=runs)
+    assert result.equals(baseline)
+    assert result.column('id').to_pylist() == _expected(fixture[1], 
runs)['id'].to_pylist()
+    assert result.column('added').null_count == len(result)
+
+
[email protected](
+    'mode', ['full', 'no_index', 'budget', 'location_budget', 'footer_bytes',
+             'footer_chunks', 'footer_items', 'cache', 'scattered'])
+def test_unsupported_or_expensive_reads_fall_back(fixture, mode):
+    path, table, file_io, counter = fixture
+    kwargs = {}
+    if mode == 'full':
+        kwargs['row_ranges'] = [(0, N - 1)]
+    elif mode == 'scattered':
+        kwargs['row_ranges'] = [(5, 6), (4500, 4540)]
+    elif mode == 'no_index':
+        pq.write_table(table, path)
+    elif mode == 'cache':
+        kwargs['row_group_cache'] = reader_module._DecodedRowGroupCache(4 * 
1024 * 1024)
+    with patch.object(page_module, '_MAX_PAGE_BYTES', 1 if mode == 'budget' 
else 32 * 1024 * 1024), \
+            patch.object(page_module, '_MAX_PAGE_LOCATIONS',
+                         1 if mode == 'location_budget' else 128 * 1024), \
+            patch.object(page_module, '_MAX_FOOTER_BYTES',
+                         1 if mode == 'footer_bytes' else 1024 * 1024), \
+            patch.object(page_module, '_MAX_FOOTER_COLUMN_CHUNKS',
+                         1 if mode == 'footer_chunks' else 1024), \
+            patch.object(page_module, '_MAX_FOOTER_ITEMS',
+                         1 if mode == 'footer_items' else 64 * 1024), \
+            patch.object(page_module.ParquetPageIndexReader, '_batches',
+                         side_effect=AssertionError('must fall back')):
+        result, _ = _read(fixture, **kwargs)
+    assert result.equals(_expected(table, kwargs.get('row_ranges', RUNS)))
+
+
[email protected]('version', ['1.0', '2.0'])
[email protected]('dictionary', [False, True])
[email protected]('compression', ['NONE', 'snappy', 'zstd', 'gzip'])
+def test_page_encodings_nulls_and_different_column_boundaries(tmp_path, 
version, dictionary, compression):
+    path = str(tmp_path / 'types.parquet')
+    count = 8192
+    table = pa.table({
+        'text': pa.array([None if i % 3 else 'v-%d' % i for i in 
range(count)]),
+        'all_null': pa.nulls(count, type=pa.int32()),
+        'flag': pa.array([i % 3 == 0 for i in range(count)]),
+        'decimal': pa.array([Decimal(i).scaleb(-2) for i in range(count)], 
pa.decimal128(20, 2)),
+        'timestamp': pa.array(range(count), pa.timestamp('ns', 
'Asia/Shanghai')),
+        'large': pa.array([bytes([i % 256]) * (i % 101) for i in 
range(count)], pa.large_binary()),
+    })
+    pq.write_table(table, path, data_page_version=version, 
use_dictionary=dictionary,
+                   compression=compression, write_page_index=True,
+                   data_page_size=1024, write_batch_size=64, 
dictionary_pagesize_limit=1024)
+    runs = [(62, 1050)]
+    with pa.OSFile(path, 'rb') as source:
+        parquet = pq.ParquetFile(source)
+        reader = page_module.ParquetPageIndexReader.create(
+            source, parquet, table.column_names, [0], 37)
+        batches = reader.read_row_group(0, runs)
+        assert batches is not None
+        assert pa.Table.from_batches(list(batches)).equals(_expected(table, 
runs))
+
+
+def test_sparse_row_indices_are_normalized(fixture):
+    rows = [N - 1, 1000, 125, 125, 126, -1, N + 1]
+    result, _ = _read(fixture, row_ranges=None, row_indices=rows)
+    assert result.equals(_expected(fixture[1], [(i, i) for i in rows]))
+    assert _read(fixture, row_ranges=[])[0] is None
+    assert _read(fixture, row_ranges=[(N, N + 10)])[0] is None
+
+
+def test_concurrent_readers_and_early_close(fixture):
+    path, table, file_io, _ = fixture
+
+    def read(_):
+        reader = reader_module.FormatPyArrowReader(
+            file_io, 'parquet', path, FIELDS, None, row_ranges=[(0, 2)], 
batch_size=71,
+            options=PAGE_INDEX_OPTIONS)
+        source = reader._parquet_source
+        try:
+            return reader.read_arrow_batch().column(0).to_pylist()
+        finally:
+            reader.close()
+            assert source.closed
+            reader.close()
+
+    with ThreadPoolExecutor(max_workers=4) as pool:
+        assert list(pool.map(read, range(8))) == [[0, 1, 2]] * 8
+
+
+def test_corrupt_offset_index_is_not_silently_ignored(fixture):
+    path = fixture[0]
+    with pa.OSFile(path, 'rb') as source:
+        parquet = pq.ParquetFile(source)
+        reader = page_module.ParquetPageIndexReader.create(source, parquet, 
['id'], [0], 71)
+        chunk = page_module._get(page_module._get(reader.footer, 4)[1][0], 
1)[1][0]
+        offset, size = page_module._get(chunk, 4), page_module._get(chunk, 5)
+        raw = source.read_at(size, offset)
+    index = page_module._Compact(raw).value(12)
+    page_module._get(index, 1)[1][0][3] = (6, 1)  # First page must start at 
row zero.
+    modified = page_module._encode(12, index)
+    assert len(modified) == size
+    with open(path, 'r+b') as output:
+        output.seek(offset)
+        output.write(modified)
+    with pytest.raises(ValueError, match='OffsetIndex row boundaries'):
+        _read(fixture, row_ranges=[(0, 2)])
+
+
+def test_offset_index_page_locations_are_bounded_before_decoding():
+    count = page_module._MAX_PAGE_LOCATIONS + 1
+    encoded = b'\x19\xfc' + page_module._unsigned(count) + b'\x00' * count + 
b'\x00'
+    with pytest.raises(page_module._PageIndexBudgetExceeded,
+                       match='page-location budget'):
+        page_module._decode_offset_index(encoded, 
page_module._MAX_PAGE_LOCATIONS)
+
+
+def test_offset_index_unknown_struct_fields_are_bounded():
+    location = {1: (6, 4), 2: (5, 1), 3: (6, 0)}
+    unknown = {field: (1, True) for field in range(1, 33)}
+    encoded = page_module._encode(
+        12, {1: (9, (12, [location])), 2: (12, unknown)})
+    with pytest.raises(page_module._PageIndexBudgetExceeded,
+                       match='object budget'):
+        page_module._decode_offset_index(encoded, 1)
+
+
+def test_page_header_unknown_struct_fields_are_bounded():
+    unknown = {field: (1, True) for field in range(1, 4097)}
+    encoded = page_module._encode(
+        12, {1: (5, 0), 2: (5, 1), 3: (5, 1), 9: (12, unknown)})
+    with pytest.raises(page_module._PageIndexBudgetExceeded,
+                       match='object budget'):
+        page_module._decode_page_header(encoded)
+
+
+def test_page_header_decoder_skips_unknown_fields_and_rejects_missing_fields():
+    encoded = page_module._encode(
+        12, {1: (5, 0), 2: (5, 11), 3: (5, 7),
+             5: (12, {1: (5, 3), 9: (9, (5, [1, 2]))}),
+             9: (12, {1: (1, True)})})
+    header, size = page_module._decode_page_header(encoded)
+    assert size == len(encoded)
+    assert page_module._get(header, 1) == 0
+    assert page_module._get(header, 2) == 11
+    assert page_module._get(header, 3) == 7
+    assert page_module._get(page_module._get(header, 5), 1) == 3
+    with pytest.raises(ValueError, match='Missing Parquet PageHeader field'):
+        page_module._decode_page_header(
+            page_module._encode(12, {1: (5, 0), 2: (5, 1)}))
+
+
+def test_page_header_budget_falls_back(fixture):
+    runs = [(4500, 4540)]
+    baseline, _ = _read(fixture, baseline=True, row_ranges=runs)
+    reader_module._reset_file_format_dataset_cache()
+    with patch.object(
+            page_module, '_decode_page_header',
+            side_effect=page_module._PageIndexBudgetExceeded('test budget')):
+        result, _ = _read(fixture, row_ranges=runs)
+    assert result.equals(baseline)
+
+
+def test_fragmented_footer_falls_back_before_generic_decoding(fixture):
+    path = fixture[0]
+    with pa.OSFile(path, 'rb') as source:
+        parquet = pq.ParquetFile(source)
+        with patch.object(page_module, '_MAX_FOOTER_COLUMN_CHUNKS', 1), \
+                patch.object(page_module._Compact, 'value',
+                             side_effect=AssertionError('must not decode 
footer')):
+            assert page_module.ParquetPageIndexReader.create(
+                source, parquet, ['id'], [0], 71) is None
+
+
+def test_wide_fallback_coalesces_offset_index_reads(tmp_path):
+    path = str(tmp_path / 'wide.parquet')
+    columns = ['column_%03d' % i for i in range(200)]
+    table = pa.table({name: range(16) for name in columns})
+    pq.write_table(table, path, write_page_index=True, use_dictionary=False,
+                   data_page_size=1024 * 1024)
+    with pa.OSFile(path, 'rb') as source:
+        reader = page_module.ParquetPageIndexReader.create(
+            source, pq.ParquetFile(source), columns, [0], 71)
+        with patch.object(page_module, '_read_exact', 
wraps=page_module._read_exact) as reads:
+            assert reader.read_row_group(0, [(0, 1)]) is None
+    assert reads.call_count == 1
+
+
+def test_index_io_errors_propagate_and_release_source(fixture):
+    path, _, file_io, _ = fixture
+    reader = reader_module.FormatPyArrowReader(
+        file_io, 'parquet', path, FIELDS, None, row_ranges=[(0, 2)], 
options=PAGE_INDEX_OPTIONS)
+    source = reader._parquet_source
+    try:
+        with patch.object(page_module, '_read_exact', 
side_effect=OSError('injected I/O failure')):
+            with pytest.raises(OSError, match='injected I/O failure'):
+                reader.read_arrow_batch()
+    finally:
+        reader.close()
+    assert source.closed
+
+
[email protected]('encoding,kind', [
+    ('DELTA_BINARY_PACKED', pa.int64()),
+    ('DELTA_LENGTH_BYTE_ARRAY', pa.string()),
+    ('DELTA_BYTE_ARRAY', pa.string()),
+    ('BYTE_STREAM_SPLIT', pa.float64()),
+])
+def test_non_dictionary_encodings(tmp_path, encoding, kind):
+    path = str(tmp_path / 'encoding.parquet')
+    values = ['common-prefix-%05d' % i for i in range(N)] if 
pa.types.is_string(kind) else range(N)
+    if pa.types.is_int64(kind):
+        values = [i ** 3 for i in range(N)]
+    table = pa.table({'value': pa.array(values, type=kind)})
+    pq.write_table(table, path, column_encoding=encoding, use_dictionary=False,
+                   write_page_index=True, data_page_size=1024, 
write_batch_size=64)
+    runs = [(8000, 8100)]
+    with pa.OSFile(path, 'rb') as source:
+        page_reader = page_module.ParquetPageIndexReader.create(
+            source, pq.ParquetFile(source), ['value'], [0], 71)
+        use_pages = page_reader.read_row_group(0, runs) is not None
+    fields = PyarrowFieldParser.to_paimon_schema(table.schema)
+    reader = reader_module.FormatPyArrowReader(
+        LocalFileIO(str(tmp_path), Options({})), 'parquet', path, fields, None,
+        row_ranges=runs, batch_size=71, options=PAGE_INDEX_OPTIONS)
+    try:
+        with patch.object(page_module.ParquetPageIndexReader, 
'_column_payload',
+                          autospec=True,
+                          
side_effect=page_module.ParquetPageIndexReader._column_payload
+                          ) as read_pages:
+            batches = []
+            while True:
+                batch = reader.read_arrow_batch()
+                if batch is None:
+                    break
+                batches.append(batch)
+        assert pa.Table.from_batches(batches).equals(table.slice(8000, 101))
+        if use_pages:
+            assert read_pages.called
+    finally:
+        reader.close()
+
+
+def test_page_header_row_count_must_agree_with_index(fixture):
+    path = fixture[0]
+    with pa.OSFile(path, 'rb') as source:
+        reader = page_module.ParquetPageIndexReader.create(source, 
pq.ParquetFile(source), ['id'], [0], 71)
+        chunk = page_module._get(page_module._get(reader.footer, 4)[1][0], 
1)[1][0]
+        offset, size = page_module._get(chunk, 4), page_module._get(chunk, 5)
+        index = page_module._Compact(source.read_at(size, offset)).value(12)
+    second = page_module._get(index, 1)[1][1]
+    second[3] = (6, page_module._get(second, 3) + 1)
+    modified = page_module._encode(12, index)
+    assert len(modified) == size
+    with open(path, 'r+b') as output:
+        output.seek(offset)
+        output.write(modified)
+    with pytest.raises(ValueError, match='page rows disagree'):
+        _read(fixture, row_ranges=[(0, 2)])
+
+
+def test_scattered_ranges_do_not_even_read_indexes(fixture):
+    with patch.object(page_module, '_read_exact', 
side_effect=AssertionError('index read')):
+        actual, _ = _read(fixture, row_ranges=[(5, 6), (4500, 4501)])
+    assert actual.equals(_expected(fixture[1], [(5, 6), (4500, 4501)]))
+
+
[email protected]('missing', [False, True])
+def test_missing_or_corrupt_parquet_still_raises(fixture, missing):
+    path = fixture[0]
+    if missing:
+        import os
+        os.remove(path)
+    else:
+        with open(path, 'wb') as output:
+            output.write(b'not a parquet file')
+    with pytest.raises((OSError, pa.ArrowInvalid)):
+        _read(fixture, row_ranges=[(0, 2)])
+
+
[email protected]('values,enabled', [
+    (None, False), ({}, True),
+    ({'parquet.filter.columnindex.enabled': 'false'}, False),
+    ({'parquet.filter.columnindex.enabled': False}, False),
+    ({'parquet.filter.columnindex.enabled': 'true'}, True),
+    ({'parquet.filter.columnindex.enabled': True}, True),
+])
+def test_page_index_switch_bypasses_metadata_processing_when_disabled(fixture, 
values, enabled):
+    options = CoreOptions(Options(values)) if values is not None else None
+    with patch.object(page_module.ParquetPageIndexReader, 'create',
+                      wraps=page_module.ParquetPageIndexReader.create) as 
create:
+        actual, _ = _read(fixture, options=options, row_ranges=[(4500, 4540)])
+    assert create.called == enabled
+    assert actual.equals(_expected(fixture[1], [(4500, 4540)]))
+
+
[email protected]('nested', [False, True])
+def test_table_option_and_copy_control_page_index_reads(tmp_path, nested):
+    from pypaimon import CatalogFactory, Schema
+
+    catalog = CatalogFactory.create({'warehouse': str(tmp_path / 'warehouse')})
+    catalog.create_database('default', False)
+    data = pa.table({'id': range(N)})
+    if nested:
+        data = data.append_column('record', pa.array([{'value': i} for i in 
range(N)]))
+    catalog.create_table('default.indexed', Schema.from_pyarrow_schema(
+        data.schema, options={
+            'row-tracking.enabled': 'true',
+            'data-evolution.enabled': 'true',
+            'parquet.filter.columnindex.enabled': 'true',
+        }), False)
+    table = catalog.get_table('default.indexed')
+    write_parquet = table.file_io.write_parquet
+
+    def write_indexed(path, arrow, **kwargs):
+        kwargs.update(write_page_index=True, data_page_size=1024, 
write_batch_size=64,
+                      use_dictionary=False)
+        return write_parquet(path, arrow, **kwargs)
+
+    builder = table.new_batch_write_builder()
+    writer, commit = builder.new_write(), builder.new_commit()
+    try:
+        with patch.object(table.file_io, 'write_parquet', 
side_effect=write_indexed):
+            writer.write_arrow(data)
+            commit.commit(writer.prepare_commit())
+    finally:
+        writer.close()
+        commit.close()
+
+    table = catalog.get_table('default.indexed')
+    for candidate, enabled in (
+            (table, True),
+            (table.copy({'parquet.filter.columnindex.enabled': 'false'}), 
False),
+            (table.copy({'parquet.filter.columnindex.enabled': 'true'}), 
True)):
+        builder = candidate.new_read_builder().with_projection(['id', 
'_ROW_ID'])
+        builder.with_filter(builder.new_predicate_builder().between('_ROW_ID', 
4500, 4540))
+        with patch.object(page_module.ParquetPageIndexReader, 'create',
+                          wraps=page_module.ParquetPageIndexReader.create) as 
create:
+            actual = 
builder.new_read().to_arrow(builder.new_scan().plan().splits())
+        assert create.called == enabled
+        assert actual.to_pydict() == {'id': list(range(4500, 4541)), 
'_ROW_ID': list(range(4500, 4541))}
+    assert table.options.parquet_column_index_enabled()
+    assert 
catalog.get_table('default.indexed').options.parquet_column_index_enabled()
+
+
[email protected]
+def nested_fixture(fixture):
+    path, original, file_io, counter = fixture
+    count = len(original)
+    child_type = pa.struct([('number', pa.int64()), ('text', pa.string())])
+    records = [None if i % 13 == 0 else {
+        'number': None if i % 11 == 0 else i,
+        'text': None if i % 7 == 0 else 
hashlib.sha256(str(i).encode()).hexdigest()
+    } for i in range(count)]
+    table = pa.table({
+        'record': pa.array(records, child_type),
+        'items': pa.array([None if i % 9 == 0 else
+                           [records[i]] * (4097 if i == N // 2 + 1 else i % 5)
+                           for i in range(count)], pa.list_(child_type)),
+        'mapping': pa.array([None if i % 9 == 0 else
+                             [('key-%d' % j, None if j == 1 else 
list(range(j)))
+                              for j in range(i % 4)] for i in range(count)],
+                            pa.map_(pa.string(), pa.list_(pa.int32()))),
+        'matrix': pa.array([None if i % 9 == 0 else
+                            [None, [], [None, i]] * (i % 3) for i in 
range(count)],
+                           pa.list_(pa.list_(pa.int64()))),
+        # Place flat columns after multiple nested physical leaves.
+        'id': original['id'],
+        'payload': original['payload'],
+    })
+    return path, table, file_io, counter
+
+
[email protected]('version', ['1.0', '2.0'])
[email protected]('dictionary', [False, True])
[email protected]('projection', ['flat', 'nested'])
+def test_nested_page_reads_preserve_structure_and_skip_bytes(
+        nested_fixture, version, dictionary, projection):
+    path, table, _, _ = nested_fixture
+    pq.write_table(table, path, write_page_index=True, 
data_page_version=version,
+                   use_dictionary=dictionary, dictionary_pagesize_limit=1024,
+                   data_page_size=2048, write_batch_size=64, row_group_size=N 
// 2)
+    names = ['payload', 'id'] if projection == 'flat' else 
list(reversed(table.column_names))
+    fields = PyarrowFieldParser.to_paimon_schema(table.select(names).schema)
+    # Cross a row-group boundary, including null parents, empty lists/maps,
+    # null elements, and multiple leaves with different page boundaries.
+    runs = [(N // 2 - 17, N // 2 + 83)]
+    baseline, baseline_reads = _read(nested_fixture, baseline=True, 
fields=fields, row_ranges=runs)
+    reader_module._reset_file_format_dataset_cache()
+    for _ in range(2):
+        with patch.object(page_module.ParquetPageIndexReader, 
'_column_payload',
+                          autospec=True, 
side_effect=page_module.ParquetPageIndexReader._column_payload
+                          ) as read_pages:
+            actual, reads = _read(nested_fixture, fields=fields, 
row_ranges=runs)
+        assert read_pages.called
+        assert actual.equals(baseline)
+        assert actual.equals(_expected(table.select(names), runs))
+        assert sum(size for _, size in reads) < sum(size for _, size in 
baseline_reads)
+
+
+def test_nested_child_projection_with_page_index(nested_fixture):
+    path, table, _, _ = nested_fixture
+    pq.write_table(table, path, write_page_index=True, use_dictionary=False,
+                   data_page_size=2048, write_batch_size=64)
+    fields = [DataField(0, 'text', AtomicType('STRING')),
+              DataField(1, 'missing', AtomicType('INT')),
+              DataField(2, 'id', AtomicType('BIGINT'))]
+    kwargs = {'fields': fields, 'nested_name_paths': [['record', 'text'], 
['record', 'absent'], ['id']],
+              'row_ranges': [(4500, 4540)]}
+    baseline, _ = _read(nested_fixture, baseline=True, **kwargs)
+    with patch.object(page_module.ParquetPageIndexReader, '_column_payload',
+                      autospec=True, 
side_effect=page_module.ParquetPageIndexReader._column_payload
+                      ) as read_pages:
+        actual, _ = _read(nested_fixture, **kwargs)
+    assert read_pages.called
+    assert actual.equals(baseline)
+    assert actual.column('text').to_pylist() == [
+        None if value is None else value['text'] for value in 
table['record'].slice(4500, 41).to_pylist()]
+    assert actual.column('missing').null_count == 41
+
+
[email protected]('version', ['1.0', '2.0'])
+def test_repeated_page_row_count_corruption_is_not_hidden(nested_fixture, 
version):
+    path, table, _, _ = nested_fixture
+    # A single-leaf nested field isolates V1's value count from its row count.
+    table = table.select(['matrix'])
+    pq.write_table(table, path, write_page_index=True, 
data_page_version=version,
+                   use_dictionary=False, data_page_size=1024, 
write_batch_size=64)
+    with pa.OSFile(path, 'rb') as source:
+        reader = page_module.ParquetPageIndexReader.create(
+            source, pq.ParquetFile(source), ['matrix'], [0], 71)
+        chunk = page_module._get(page_module._get(reader.footer, 4)[1][0], 
1)[1][0]
+        offset, size = page_module._get(chunk, 4), page_module._get(chunk, 5)
+        index = page_module._Compact(source.read_at(size, offset)).value(12)
+    second = page_module._get(index, 1)[1][1]
+    second[3] = (6, page_module._get(second, 3) + 1)
+    modified = page_module._encode(12, index)
+    assert len(modified) == size
+    with open(path, 'r+b') as output:
+        output.seek(offset)
+        output.write(modified)
+    fields = PyarrowFieldParser.to_paimon_schema(table.schema)
+    with pytest.raises((ValueError, pa.ArrowInvalid), match='rows|row'):
+        _read(nested_fixture, fields=fields, row_ranges=[(0, 2)])
+
+
+def 
test_nested_alignment_can_fall_back_when_no_pages_can_be_skipped(nested_fixture):
+    path, table, _, _ = nested_fixture
+    # One leaf has a single page, forcing the field's common span to the full 
group.
+    table = pa.table({'record': pa.StructArray.from_arrays(
+        [pa.array([True] * N), table['payload'].combine_chunks()], 
names=['flag', 'text'])})
+    pq.write_table(table, path, write_page_index=True, use_dictionary=False,
+                   data_page_size=4096, write_batch_size=64)
+    fields = PyarrowFieldParser.to_paimon_schema(table.schema)
+    with patch.object(page_module.ParquetPageIndexReader, '_column_payload',
+                      side_effect=AssertionError('must fall back before 
reading pages')):
+        actual, _ = _read(nested_fixture, fields=fields, row_ranges=[(4500, 
4540)])
+    assert actual.equals(table.slice(4500, 41))

Reply via email to