From 9932c187addc546d702a885ffcd78d634ac3549b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Fri, 18 Sep 2026 02:23:09 -0700 Subject: [PATCH 01/10] [python] Add opt-in Parquet OffsetIndex reads for row windows --- paimon-python/README.md | 22 ++ .../pypaimon/common/options/core_options.py | 15 + .../read/reader/format_pyarrow_reader.py | 66 +++- .../read/reader/parquet_page_index_reader.py | 368 ++++++++++++++++++ .../pypaimon/tests/parquet_page_index_test.py | 336 ++++++++++++++++ 5 files changed, 805 insertions(+), 2 deletions(-) create mode 100644 paimon-python/pypaimon/read/reader/parquet_page_index_reader.py create mode 100644 paimon-python/pypaimon/tests/parquet_page_index_test.py diff --git a/paimon-python/README.md b/paimon-python/README.md index b3541b1bdb48..7ffc4fc0f26a 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -105,6 +105,28 @@ pip3 install dist/*.tar.gz The command will install the package and core dependencies to your local Python environment. +# Parquet page-index reads + +`read.parquet.page-index.enabled` is disabled by default. Enable it for a read +using a table copy, without changing persisted table options: + +```python +indexed_table = table.copy({"read.parquet.page-index.enabled": "true"}) +builder = indexed_table.new_read_builder() +# On a row-tracking table, select a contiguous row-ID window. +builder.with_filter(builder.new_predicate_builder().between("_ROW_ID", 100, 199)) +rows = builder.new_read().to_arrow(builder.new_scan().plan().splits()) +``` + +The optimization uses existing Parquet OffsetIndexes for one contiguous window +per row group in flat schemas. Disjoint windows, nested schemas, missing indexes, +older Arrow versions without index metadata support, and decoded row-group cache +reads retain the ordinary reader. Large or unprofitable page selections also fall +back. Enabling this option does not enable ColumnIndex predicate filtering or +force unsupported reads through the page-index path. It can reduce bytes read; +fewer OSS HEAD/GET requests or lower latency are not guaranteed. Set the option to +`"false"` to bypass page-index processing entirely. + # 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 7df8ee0ef25f..4293c0835be0 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1128,6 +1128,18 @@ class CoreOptions: .with_description("Read batch size for any file format if it supports.") ) + READ_PARQUET_PAGE_INDEX_ENABLED: ConfigOption[bool] = ( + ConfigOptions.key("read.parquet.page-index.enabled") + .boolean_type() + .default_value(False) + .with_description( + "Enable PyPaimon Parquet OffsetIndex reads for contiguous row windows. " + "Disabled by default. Requires flat schemas and existing offset indexes; " + "unsupported or expensive selections use the ordinary reader. " + "Does not enable ColumnIndex predicate filtering." + ) + ) + READ_PARALLELISM: ConfigOption[int] = ( ConfigOptions.key("read.parallelism") .int_type() @@ -1878,6 +1890,9 @@ def local_cache_whitelist(self) -> str: def read_batch_size(self, default=None) -> int: return self.options.get(CoreOptions.READ_BATCH_SIZE, default or 1024) + def read_parquet_page_index_enabled(self) -> bool: + return self.options.get(CoreOptions.READ_PARQUET_PAGE_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 33999fec75f8..aa96d29eb827 100644 --- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py +++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py @@ -518,6 +518,8 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, # 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 @@ -526,8 +528,26 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, 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.read_parquet_page_index_enabled() + and self._row_group_cache is None + and not self._has_nested_path + 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( @@ -535,6 +555,11 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, 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: @@ -611,6 +636,36 @@ def _iter_row_group_batches(self): if out.num_rows: yield out + def _iter_page_index_batches(self, selected_infos, runs): + 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 self._select_existing_fields(batch) + else: + try: + for batch in batches: + yield self._select_existing_fields(batch) + finally: + batches.close() + def _read_parquet_row_group_batches(self, row_group, columns): return self._parquet_file.iter_batches( row_groups=[row_group], @@ -852,12 +907,19 @@ def _cast_orc_time_columns(self, batch): 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 000000000000..06141c18d208 --- /dev/null +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -0,0 +1,368 @@ +# 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 windows of flat Parquet columns 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 schemas and files without indexes use the +normal reader. Disjoint ranges use the normal reader too: 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 data retained by the temporary column files, independently of +# the number of requested rows and the size of the source row group. +_MAX_PAGE_BYTES = 32 * 1024 * 1024 +_MAX_INDEX_BYTES = 8 * 1024 * 1024 + + +class _Compact: + """Thrift compact values used by Parquet metadata (no generated bindings).""" + + def __init__(self, data): + self.data = memoryview(data) + self.position = 0 + + 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") + 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") + fields[field] = field_kind, self.value(field_kind, depth + 1) + previous = field + raise ValueError("Unsupported Parquet compact type: {}".format(kind)) + + +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 + + +class ParquetPageIndexReader: + def __init__(self, source, metadata, schema, footer, columns, batch_size): + self.source = source + self.metadata = metadata + self.schema = schema + self.footer = footer + self.columns = columns + 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 + # Flat columns have one schema element and one physical column each. + if not columns or any(pa.types.is_nested(field.type) for field in schema): + return None + if len(schema) != metadata.num_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(all(getattr(metadata.row_group(group).column(index), + "has_offset_index", False) for index in indices) + for group in row_groups): + return None + output = pa.BufferOutputStream() + metadata.write_metadata_file(output) + serialized = output.getvalue().to_pybytes() + length = struct.unpack("= row_count: + return None + chunks = _get(row_group, 1)[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 self.columns): + return None + plans = [] + selected_bytes = 0 + index_bytes = 0 + full_bytes = 0 + for index in sorted(self.columns): + chunk = chunks[index] + index_size = _get(chunk, 5) + if index_size > _MAX_INDEX_BYTES: + return None + raw = _read_exact(self.source, _get(chunk, 4), index_size) + index_bytes += index_size + locations = _get(_Compact(raw).value(12), 1)[1] + column = _get(chunk, 3) + data_offset = _get(column, 9) + dictionary_offset = _get(column, 11, data_offset) + chunk_end = dictionary_offset + _get(column, 7) + starts = [_get(page, 3) 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 = (_get(page, field) for field in (1, 2, 3)) + 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 _get(locations[0], 1) != data_offset or dictionary_offset > data_offset: + raise ValueError("Invalid Parquet OffsetIndex first page") + lower, upper = runs[0] + selected = range(bisect.bisect_right(starts, lower) - 1, + bisect.bisect_right(starts, upper)) + pages = [] + infos = [] + for position in selected: + page = locations[position] + pages.append((_get(page, 1), _get(page, 2))) + end = starts[position + 1] if position + 1 < len(starts) else row_count + infos.append((starts[position], end - starts[position])) + dictionary_size = data_offset - dictionary_offset + selected_bytes += dictionary_size + sum(size for _, size in pages) + full_bytes += _get(column, 7) + plans.append((index, column, dictionary_offset, dictionary_size, pages, infos)) + if (selected_bytes + index_bytes >= full_bytes + or selected_bytes > _MAX_PAGE_BYTES): + return None + return self._batches(plans, runs) + + def _column_batches(self, plan, runs): + from pypaimon.read.reader.format_pyarrow_reader import _RowRunSlicer + + 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 + for position, (_, length) in enumerate(ranges): + parser = _Compact(memoryview(payload)[cursor:cursor + length]) + header = parser.value(12) + if parser.position + _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: + actual = _get(_get(header, 5), 1) + elif page_type == 3: + page_header = _get(header, 8) + actual = _get(page_header, 3) + if _get(page_header, 1) != actual: + raise ValueError("Invalid flat Parquet data page") + else: + raise ValueError("Invalid Parquet data page type") + if actual != expected: + raise ValueError("Parquet page rows disagree with OffsetIndex") + uncompressed_size += parser.position + _get(header, 2) + cursor += length + + num_rows = sum(count for _, count in infos) + patched_column = {key: value for key, value in column.items() if key <= 8} + patched_column.update({5: (6, num_rows), 6: (6, uncompressed_size), + 7: (6, len(payload)), 9: (6, 4 + dictionary_size)}) + if dictionary_size: + patched_column[11] = (6, 4) + patched_group = {1: (9, (12, [{2: (6, 0), 3: (12, patched_column)}])), + 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, elements[index + 1]])), + 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"PAR1" + payload + encoded + struct.pack(" Date: Fri, 18 Sep 2026 02:29:49 -0700 Subject: [PATCH 02/10] [python] Reuse the Java Parquet page-index configuration key --- paimon-python/README.md | 7 ++++--- .../pypaimon/common/options/core_options.py | 11 ++++++----- .../read/reader/format_pyarrow_reader.py | 2 +- .../pypaimon/tests/parquet_page_index_test.py | 16 ++++++++-------- 4 files changed, 19 insertions(+), 17 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index 7ffc4fc0f26a..27d084fcd367 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -107,11 +107,12 @@ The command will install the package and core dependencies to your local Python # Parquet page-index reads -`read.parquet.page-index.enabled` is disabled by default. Enable it for a read -using a table copy, without changing persisted table options: +PyPaimon uses the same `parquet.filter.columnindex.enabled` option as the Java +reader. It defaults to `false` in Python and `true` in Java. Enable it for a Python +read using a table copy, without changing persisted table options: ```python -indexed_table = table.copy({"read.parquet.page-index.enabled": "true"}) +indexed_table = table.copy({"parquet.filter.columnindex.enabled": "true"}) builder = indexed_table.new_read_builder() # On a row-tracking table, select a contiguous row-ID window. builder.with_filter(builder.new_predicate_builder().between("_ROW_ID", 100, 199)) diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index 4293c0835be0..f9964802fc68 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1128,13 +1128,14 @@ class CoreOptions: .with_description("Read batch size for any file format if it supports.") ) - READ_PARQUET_PAGE_INDEX_ENABLED: ConfigOption[bool] = ( - ConfigOptions.key("read.parquet.page-index.enabled") + PARQUET_COLUMN_INDEX_ENABLED: ConfigOption[bool] = ( + ConfigOptions.key("parquet.filter.columnindex.enabled") .boolean_type() .default_value(False) .with_description( "Enable PyPaimon Parquet OffsetIndex reads for contiguous row windows. " - "Disabled by default. Requires flat schemas and existing offset indexes; " + "Uses the same key as Java, but defaults to false in Python (true in Java). " + "Requires flat schemas and existing offset indexes; " "unsupported or expensive selections use the ordinary reader. " "Does not enable ColumnIndex predicate filtering." ) @@ -1890,8 +1891,8 @@ def local_cache_whitelist(self) -> str: def read_batch_size(self, default=None) -> int: return self.options.get(CoreOptions.READ_BATCH_SIZE, default or 1024) - def read_parquet_page_index_enabled(self) -> bool: - return self.options.get(CoreOptions.READ_PARQUET_PAGE_INDEX_ENABLED) + 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 aa96d29eb827..2523bb848819 100644 --- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py +++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py @@ -534,7 +534,7 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, self._parquet_file = pq.ParquetFile(self._parquet_source) if (self._selected_parquet_row_groups is not None and options is not None - and options.read_parquet_page_index_enabled() + and options.parquet_column_index_enabled() and self._row_group_cache is None and not self._has_nested_path and not self._bounded_variant_read): diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 8aeda2d82873..534b0c0fd342 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -41,7 +41,7 @@ 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({'read.parquet.page-index.enabled': 'true'})) +PAGE_INDEX_OPTIONS = CoreOptions(Options({'parquet.filter.columnindex.enabled': 'true'})) @pytest.fixture @@ -282,10 +282,10 @@ def test_missing_or_corrupt_parquet_still_raises(fixture, missing): @pytest.mark.parametrize('values,enabled', [ (None, False), ({}, False), - ({'read.parquet.page-index.enabled': 'false'}, False), - ({'read.parquet.page-index.enabled': False}, False), - ({'read.parquet.page-index.enabled': 'true'}, True), - ({'read.parquet.page-index.enabled': True}, 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 @@ -324,7 +324,7 @@ def write_indexed(path, arrow, **kwargs): table = catalog.get_table('default.indexed') for value in ('true', 'false', 'true'): - copied = table.copy({'read.parquet.page-index.enabled': value}) + copied = table.copy({'parquet.filter.columnindex.enabled': value}) builder = copied.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', @@ -332,5 +332,5 @@ def write_indexed(path, arrow, **kwargs): actual = builder.new_read().to_arrow(builder.new_scan().plan().splits()) assert create.called == (value == 'true') assert actual.to_pydict() == {'id': list(range(4500, 4541)), '_ROW_ID': list(range(4500, 4541))} - assert not table.options.read_parquet_page_index_enabled() - assert not catalog.get_table('default.indexed').options.read_parquet_page_index_enabled() + assert not table.options.parquet_column_index_enabled() + assert not catalog.get_table('default.indexed').options.parquet_column_index_enabled() From 804288cf9340d1faae35df46480609442a67012b Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Fri, 18 Sep 2026 03:18:52 -0700 Subject: [PATCH 03/10] [python] Support nested fields in Parquet OffsetIndex window reads --- paimon-python/README.md | 12 +- .../pypaimon/common/options/core_options.py | 4 +- .../read/reader/format_pyarrow_reader.py | 6 +- .../read/reader/parquet_page_index_reader.py | 156 +++++++++++++----- .../pypaimon/tests/parquet_page_index_test.py | 127 +++++++++++++- 5 files changed, 248 insertions(+), 57 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index 27d084fcd367..89c977a7d3e6 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -120,12 +120,18 @@ rows = builder.new_read().to_arrow(builder.new_scan().plan().splits()) ``` The optimization uses existing Parquet OffsetIndexes for one contiguous window -per row group in flat schemas. Disjoint windows, nested schemas, missing indexes, +per row group, including STRUCT, ARRAY, and MAP fields. Nested leaf columns are +aligned to common page row boundaries while retaining the original schema and +null/empty collection semantics. Unselected nested fields do not disable reads +of ordinary columns. Disjoint windows, missing indexes, older Arrow versions without index metadata support, and decoded row-group cache reads retain the ordinary reader. Large or unprofitable page selections also fall -back. Enabling this option does not enable ColumnIndex predicate filtering or +back, including nested fields whose common boundaries require reading the whole +leaf chunks without savings. VARIANT reads retain their existing path. +Enabling this option does not enable ColumnIndex predicate filtering or force unsupported reads through the page-index path. It can reduce bytes read; -fewer OSS HEAD/GET requests or lower latency are not guaranteed. Set the option to +page seeks can also increase OSS GET requests despite reducing bytes, so lower +latency or request cost is not guaranteed. Set the option to `"false"` to bypass page-index processing entirely. # Native scan planning diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index f9964802fc68..301d07f60cc7 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1135,8 +1135,8 @@ class CoreOptions: .with_description( "Enable PyPaimon Parquet OffsetIndex reads for contiguous row windows. " "Uses the same key as Java, but defaults to false in Python (true in Java). " - "Requires flat schemas and existing offset indexes; " - "unsupported or expensive selections use the ordinary reader. " + "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." ) ) diff --git a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py index 2523bb848819..c37d76d7c3d5 100644 --- a/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py +++ b/paimon-python/pypaimon/read/reader/format_pyarrow_reader.py @@ -536,7 +536,6 @@ def __init__(self, file_io: FileIO, file_format: str, file_path: str, and options is not None and options.parquet_column_index_enabled() and self._row_group_cache is None - and not self._has_nested_path and not self._bounded_variant_read): from pypaimon.read.reader.parquet_page_index_reader import ( ParquetPageIndexReader, @@ -637,6 +636,7 @@ def _iter_row_group_batches(self): 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): @@ -658,11 +658,11 @@ def _iter_page_index_batches(self, selected_infos, runs): batch = slicer.next_batch(raw) if batch is None: break - yield self._select_existing_fields(batch) + yield select(batch) else: try: for batch in batches: - yield self._select_existing_fields(batch) + yield select(batch) finally: batches.close() diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py index 06141c18d208..acb6608ae562 100644 --- a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -15,13 +15,14 @@ # specific language governing permissions and limitations # under the License. -"""Read contiguous windows of flat Parquet columns using their OffsetIndex. +"""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 schemas and files without indexes use the -normal reader. Disjoint ranges use the normal reader too: per-page seeks can -amplify requests in filesystems that prefetch remote data (including Jindo). +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 @@ -162,29 +163,27 @@ def _read_exact(source, offset, length): class ParquetPageIndexReader: - def __init__(self, source, metadata, schema, footer, columns, batch_size): + 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 - # Flat columns have one schema element and one physical column each. - if not columns or any(pa.types.is_nested(field.type) for field in schema): - return None - if len(schema) != metadata.num_columns or len(set(schema.names)) != len(schema): + 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(all(getattr(metadata.row_group(group).column(index), - "has_offset_index", False) for index in indices) - for group in row_groups): + 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 output = pa.BufferOutputStream() metadata.write_metadata_file(output) @@ -194,9 +193,32 @@ def create(cls, source, parquet_file, columns, row_groups, batch_size): if 8 in footer or 9 in footer: return None # Encrypted pages need the original file identity/AAD. elements = _get(footer, 2)[1] - if len(elements) != metadata.num_columns + 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, batch_size) + 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.""" @@ -210,16 +232,18 @@ def read_row_group(self, group, runs): 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 self.columns): + for index in physical_columns): return None + indexed = {} plans = [] selected_bytes = 0 index_bytes = 0 full_bytes = 0 - for index in sorted(self.columns): + for index in sorted(physical_columns): chunk = chunks[index] index_size = _get(chunk, 5) if index_size > _MAX_INDEX_BYTES: @@ -244,28 +268,42 @@ def read_row_group(self, group, runs): previous_end, previous_row = offset + size, first_row if _get(locations[0], 1) != data_offset or dictionary_offset > data_offset: raise ValueError("Invalid Parquet OffsetIndex first page") - lower, upper = runs[0] - selected = range(bisect.bisect_right(starts, lower) - 1, - bisect.bisect_right(starts, upper)) - pages = [] - infos = [] - for position in selected: - page = locations[position] - pages.append((_get(page, 1), _get(page, 2))) - end = starts[position + 1] if position + 1 < len(starts) else row_count - infos.append((starts[position], end - starts[position])) - dictionary_size = data_offset - dictionary_offset - selected_bytes += dictionary_size + sum(size for _, size in pages) + indexed[index] = (column, dictionary_offset, data_offset - dictionary_offset, + locations, starts) full_bytes += _get(column, 7) - plans.append((index, column, dictionary_offset, dictionary_size, pages, infos)) + 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((_get(page, 1), _get(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 return self._batches(plans, runs) - def _column_batches(self, plan, runs): - from pypaimon.read.reader.format_pyarrow_reader import _RowRunSlicer - + 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. @@ -279,6 +317,8 @@ def _column_batches(self, plan, runs): 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): parser = _Compact(memoryview(payload)[cursor:cursor + length]) header = parser.value(12) @@ -291,43 +331,73 @@ def _column_batches(self, plan, runs): expected = infos[position - bool(dictionary_size)][1] page_type = _get(header, 1) if page_type == 0: - actual = _get(_get(header, 5), 1) + 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) - if _get(page_header, 1) != actual: - raise ValueError("Invalid flat Parquet data page") + 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: + if actual != expected or values < expected: raise ValueError("Parquet page rows disagree with OffsetIndex") + num_values += values uncompressed_size += parser.position + _get(header, 2) cursor += length - num_rows = sum(count for _, count in infos) patched_column = {key: value for key, value in column.items() if key <= 8} - patched_column.update({5: (6, num_rows), 6: (6, uncompressed_size), + 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) - patched_group = {1: (9, (12, [{2: (6, 0), 3: (12, patched_column)}])), + 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, elements[index + 1]])), + 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"PAR1" + payload + encoded + struct.pack(" 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) diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 534b0c0fd342..082f6ad953ed 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -30,7 +30,7 @@ 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 +from pypaimon.schema.data_types import AtomicType, DataField, PyarrowFieldParser from pypaimon.tests.parquet_metadata_cache_test import _CountingLocalFileSystem @@ -113,7 +113,7 @@ def test_ranges_projection_missing_fields_and_fallback_in_same_file(fixture): assert result.column('added').null_count == len(result) -@pytest.mark.parametrize('mode', ['full', 'no_index', 'nested', 'budget', 'cache', 'scattered']) +@pytest.mark.parametrize('mode', ['full', 'no_index', 'budget', 'cache', 'scattered']) def test_unsupported_or_expensive_reads_fall_back(fixture, mode): path, table, file_io, counter = fixture kwargs = {} @@ -123,9 +123,6 @@ def test_unsupported_or_expensive_reads_fall_back(fixture, mode): kwargs['row_ranges'] = [(5, 6), (4500, 4540)] elif mode == 'no_index': pq.write_table(table, path) - elif mode == 'nested': - pq.write_table(table.append_column('nested', pa.array([[i] for i in range(N)])), - path, write_page_index=True) 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), \ @@ -296,12 +293,15 @@ def test_page_index_switch_bypasses_metadata_processing_when_disabled(fixture, v assert actual.equals(_expected(fixture[1], [(4500, 4540)])) -def test_table_copy_can_enable_and_disable_page_index_reads(tmp_path): +@pytest.mark.parametrize('nested', [False, True]) +def test_table_copy_can_enable_and_disable_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'}), False) table = catalog.get_table('default.indexed') @@ -334,3 +334,118 @@ def write_indexed(path, arrow, **kwargs): assert actual.to_pydict() == {'id': list(range(4500, 4541)), '_ROW_ID': list(range(4500, 4541))} assert not table.options.parquet_column_index_enabled() assert not catalog.get_table('default.indexed').options.parquet_column_index_enabled() + + +@pytest.fixture +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 + + +@pytest.mark.parametrize('version', ['1.0', '2.0']) +@pytest.mark.parametrize('dictionary', [False, True]) +@pytest.mark.parametrize('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 + + +@pytest.mark.parametrize('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)) From c07c0798990b99063f2c831d320ff5ae024c8538 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Fri, 18 Sep 2026 19:21:43 -0700 Subject: [PATCH 04/10] [python] Clarify Parquet page-index option semantics --- paimon-python/README.md | 26 ++++--------------- .../pypaimon/common/options/core_options.py | 4 +-- .../pypaimon/tests/parquet_page_index_test.py | 22 ++++++++++------ 3 files changed, 21 insertions(+), 31 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index 89c977a7d3e6..cd4a066be15c 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -107,32 +107,16 @@ The command will install the package and core dependencies to your local Python # Parquet page-index reads -PyPaimon uses the same `parquet.filter.columnindex.enabled` option as the Java -reader. It defaults to `false` in Python and `true` in Java. Enable it for a Python -read using a table copy, without changing persisted table options: +For row-tracking tables with a Parquet OffsetIndex, PyPaimon can read a +contiguous `_ROW_ID` range without decoding the full row group. Enable it +with the table option: ```python indexed_table = table.copy({"parquet.filter.columnindex.enabled": "true"}) -builder = indexed_table.new_read_builder() -# On a row-tracking table, select a contiguous row-ID window. -builder.with_filter(builder.new_predicate_builder().between("_ROW_ID", 100, 199)) -rows = builder.new_read().to_arrow(builder.new_scan().plan().splits()) ``` -The optimization uses existing Parquet OffsetIndexes for one contiguous window -per row group, including STRUCT, ARRAY, and MAP fields. Nested leaf columns are -aligned to common page row boundaries while retaining the original schema and -null/empty collection semantics. Unselected nested fields do not disable reads -of ordinary columns. Disjoint windows, missing indexes, -older Arrow versions without index metadata support, and decoded row-group cache -reads retain the ordinary reader. Large or unprofitable page selections also fall -back, including nested fields whose common boundaries require reading the whole -leaf chunks without savings. VARIANT reads retain their existing path. -Enabling this option does not enable ColumnIndex predicate filtering or -force unsupported reads through the page-index path. It can reduce bytes read; -page seeks can also increase OSS GET requests despite reducing bytes, so lower -latency or request cost is not guaranteed. Set the option to -`"false"` to bypass page-index processing entirely. +Unsupported reads use the normal path. Reading fewer bytes may require more +object-store requests. # Native scan planning diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index 301d07f60cc7..7521bd7a0414 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1133,8 +1133,8 @@ class CoreOptions: .boolean_type() .default_value(False) .with_description( - "Enable PyPaimon Parquet OffsetIndex reads for contiguous row windows. " - "Uses the same key as Java, but defaults to false in Python (true in Java). " + "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." diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 082f6ad953ed..87f70d041b8d 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -294,7 +294,7 @@ def test_page_index_switch_bypasses_metadata_processing_when_disabled(fixture, v @pytest.mark.parametrize('nested', [False, True]) -def test_table_copy_can_enable_and_disable_page_index_reads(tmp_path, nested): +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')}) @@ -303,7 +303,11 @@ def test_table_copy_can_enable_and_disable_page_index_reads(tmp_path, nested): 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'}), False) + 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 @@ -323,17 +327,19 @@ def write_indexed(path, arrow, **kwargs): commit.close() table = catalog.get_table('default.indexed') - for value in ('true', 'false', 'true'): - copied = table.copy({'parquet.filter.columnindex.enabled': value}) - builder = copied.new_read_builder().with_projection(['id', '_ROW_ID']) + 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 == (value == 'true') + assert create.called == enabled assert actual.to_pydict() == {'id': list(range(4500, 4541)), '_ROW_ID': list(range(4500, 4541))} - assert not table.options.parquet_column_index_enabled() - assert not catalog.get_table('default.indexed').options.parquet_column_index_enabled() + assert table.options.parquet_column_index_enabled() + assert catalog.get_table('default.indexed').options.parquet_column_index_enabled() @pytest.fixture From dea625231c595f1cec9c1fffabed966c2c07db7a Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Fri, 18 Sep 2026 19:38:57 -0700 Subject: [PATCH 05/10] [python] Allow Parquet page-index encoding fallback --- .../pypaimon/tests/parquet_page_index_test.py | 26 ++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 87f70d041b8d..971165db5470 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -233,11 +233,31 @@ def test_non_dictionary_encodings(tmp_path, encoding, kind): 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: - reader = page_module.ParquetPageIndexReader.create( + page_reader = page_module.ParquetPageIndexReader.create( source, pq.ParquetFile(source), ['value'], [0], 71) - actual = pa.Table.from_batches(list(reader.read_row_group(0, [(8000, 8100)]))) - assert actual.equals(table.slice(8000, 101)) + 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): From d335926502a7503640258e166b4a3cfe2bf35ec8 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sat, 19 Sep 2026 19:28:08 -0700 Subject: [PATCH 06/10] [python] Align Parquet page-index option default --- paimon-python/README.md | 6 +++--- paimon-python/pypaimon/common/options/core_options.py | 2 +- paimon-python/pypaimon/tests/parquet_page_index_test.py | 2 +- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/paimon-python/README.md b/paimon-python/README.md index cd4a066be15c..143ebcea0cc8 100644 --- a/paimon-python/README.md +++ b/paimon-python/README.md @@ -108,11 +108,11 @@ The command will install the package and core dependencies to your local Python # 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. Enable it -with the table option: +contiguous `_ROW_ID` range without decoding the full row group. This is enabled +by default and can be disabled with the table option: ```python -indexed_table = table.copy({"parquet.filter.columnindex.enabled": "true"}) +table = table.copy({"parquet.filter.columnindex.enabled": "false"}) ``` Unsupported reads use the normal path. Reading fewer bytes may require more diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index 7521bd7a0414..181ffc2190d2 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -1131,7 +1131,7 @@ class CoreOptions: PARQUET_COLUMN_INDEX_ENABLED: ConfigOption[bool] = ( ConfigOptions.key("parquet.filter.columnindex.enabled") .boolean_type() - .default_value(False) + .default_value(True) .with_description( "Enable Parquet page-index pruning. PyPaimon currently uses OffsetIndex " "metadata for contiguous row windows. " diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 971165db5470..f7e7dc27df63 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -298,7 +298,7 @@ def test_missing_or_corrupt_parquet_still_raises(fixture, missing): @pytest.mark.parametrize('values,enabled', [ - (None, False), ({}, False), + (None, False), ({}, True), ({'parquet.filter.columnindex.enabled': 'false'}, False), ({'parquet.filter.columnindex.enabled': False}, False), ({'parquet.filter.columnindex.enabled': 'true'}, True), From 14cc801e0473e862790a94a0a3b4a08d18b35d51 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 20 Sep 2026 00:13:33 -0700 Subject: [PATCH 07/10] [python] Bound and batch Parquet OffsetIndex reads --- .../read/reader/parquet_page_index_reader.py | 201 ++++++++++++++++-- .../pypaimon/tests/parquet_page_index_test.py | 21 ++ 2 files changed, 210 insertions(+), 12 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py index acb6608ae562..2ec7fb70ed51 100644 --- a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -33,10 +33,11 @@ import pyarrow.parquet as pq -# Bound encoded data retained by the temporary column files, independently of -# the number of requested rows and the size of the source row group. +# Bound encoded page data retained by the temporary column files. _MAX_PAGE_BYTES = 32 * 1024 * 1024 +# Bound all serialized OffsetIndexes and their retained typed PageLocations. _MAX_INDEX_BYTES = 8 * 1024 * 1024 +_MAX_PAGE_LOCATIONS = 128 * 1024 class _Compact: @@ -104,6 +105,156 @@ def value(self, kind, depth=0): raise ValueError("Unsupported Parquet compact type: {}".format(kind)) +class _OffsetIndexDecoder: + """Decode typed PageLocations without materializing generic Thrift trees.""" + + def __init__(self, data, max_locations): + self.parser = _Compact(data) + self.max_locations = max_locations + self.remaining_items = max_locations * 4 + + def decode(self): + locations = None + previous = 0 + seen = set() + while True: + field = self._field(previous) + if field is None: + break + field_id, kind = field + if field_id in seen: + raise ValueError("Duplicate Parquet OffsetIndex field") + seen.add(field_id) + previous = field_id + if field_id == 1: + 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 _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") + 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 _locations(self): + count, element = self._collection() + if element != 12 or count > self.max_locations: + raise ValueError("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 + seen = set() + while True: + field = self._field(previous) + if field is None: + break + field_id, kind = field + if field_id in seen: + raise ValueError("Duplicate Parquet PageLocation field") + seen.add(field_id) + previous = field_id + if 1 <= field_id <= 3: + 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 + + def _consume(self, count): + if count > self.remaining_items: + raise ValueError("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 + seen = set() + while True: + field = self._field(previous) + if field is None: + return + field_id, field_kind = field + if field_id in seen: + raise ValueError("Duplicate Parquet compact field") + seen.add(field_id) + previous = field_id + self._skip(field_kind, depth + 1) + raise ValueError("Unsupported Parquet compact type: {}".format(kind)) + + +def _decode_offset_index(data, max_locations): + return _OffsetIndexDecoder(data, max_locations).decode() + + def _unsigned(value): result = bytearray() while value >= 128: @@ -162,6 +313,27 @@ def _read_exact(source, offset, length): 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 @@ -241,32 +413,37 @@ def read_row_group(self, group, runs): indexed = {} plans = [] selected_bytes = 0 - index_bytes = 0 full_bytes = 0 + index_ranges = [] for index in sorted(physical_columns): chunk = chunks[index] index_size = _get(chunk, 5) - if index_size > _MAX_INDEX_BYTES: - return None - raw = _read_exact(self.source, _get(chunk, 4), index_size) - index_bytes += index_size - locations = _get(_Compact(raw).value(12), 1)[1] + 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] + locations = _decode_offset_index(raw_indexes[index], remaining_locations) + 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 = [_get(page, 3) for page in locations] + 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 = (_get(page, field) for field in (1, 2, 3)) + 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 _get(locations[0], 1) != data_offset or dictionary_offset > data_offset: + 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) @@ -292,7 +469,7 @@ def read_row_group(self, group, runs): pages, infos = [], [] for position in selected: page = locations[position] - pages.append((_get(page, 1), _get(page, 2))) + 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) diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index f7e7dc27df63..a619b75e11b3 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -205,6 +205,27 @@ def test_corrupt_offset_index_is_not_silently_ignored(fixture): _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(ValueError, match='page-location budget'): + page_module._decode_offset_index(encoded, page_module._MAX_PAGE_LOCATIONS) + + +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( From 09b2078b3e6b62416437685210fe175774f05c36 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 20 Sep 2026 00:33:33 -0700 Subject: [PATCH 08/10] [python] Fall back when OffsetIndex budgets are exceeded --- .../read/reader/parquet_page_index_reader.py | 19 +++++++++++++++---- .../pypaimon/tests/parquet_page_index_test.py | 8 ++++++-- 2 files changed, 21 insertions(+), 6 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py index 2ec7fb70ed51..628e6930b4f9 100644 --- a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -40,6 +40,10 @@ _MAX_PAGE_LOCATIONS = 128 * 1024 +class _PageIndexBudgetExceeded(Exception): + pass + + class _Compact: """Thrift compact values used by Parquet metadata (no generated bindings).""" @@ -157,8 +161,11 @@ def _collection(self): def _locations(self): count, element = self._collection() - if element != 12 or count > self.max_locations: - raise ValueError("Parquet OffsetIndex exceeds page-location budget") + 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)] @@ -191,7 +198,8 @@ def _location(self): def _consume(self, count): if count > self.remaining_items: - raise ValueError("Parquet compact metadata exceeds object budget") + raise _PageIndexBudgetExceeded( + "Parquet compact metadata exceeds object budget") self.remaining_items -= count def _skip_collection_value(self, kind, depth): @@ -426,7 +434,10 @@ def read_row_group(self, group, runs): remaining_locations = _MAX_PAGE_LOCATIONS for index in sorted(physical_columns): chunk = chunks[index] - locations = _decode_offset_index(raw_indexes[index], remaining_locations) + 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) diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index a619b75e11b3..5b49e2c7a95c 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -113,7 +113,8 @@ def test_ranges_projection_missing_fields_and_fallback_in_same_file(fixture): assert result.column('added').null_count == len(result) -@pytest.mark.parametrize('mode', ['full', 'no_index', 'budget', 'cache', 'scattered']) +@pytest.mark.parametrize( + 'mode', ['full', 'no_index', 'budget', 'location_budget', 'cache', 'scattered']) def test_unsupported_or_expensive_reads_fall_back(fixture, mode): path, table, file_io, counter = fixture kwargs = {} @@ -126,6 +127,8 @@ def test_unsupported_or_expensive_reads_fall_back(fixture, mode): 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.ParquetPageIndexReader, '_batches', side_effect=AssertionError('must fall back')): result, _ = _read(fixture, **kwargs) @@ -208,7 +211,8 @@ def test_corrupt_offset_index_is_not_silently_ignored(fixture): 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(ValueError, match='page-location budget'): + with pytest.raises(page_module._PageIndexBudgetExceeded, + match='page-location budget'): page_module._decode_offset_index(encoded, page_module._MAX_PAGE_LOCATIONS) From 023690452384e4504140b740fd514ab7f33ff382 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 20 Sep 2026 01:13:27 -0700 Subject: [PATCH 09/10] [python] Bound Parquet page index metadata decoding --- .../read/reader/parquet_page_index_reader.py | 30 +++++++++++-------- .../pypaimon/tests/parquet_page_index_test.py | 28 ++++++++++++++++- 2 files changed, 44 insertions(+), 14 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py index 628e6930b4f9..bdf6283cbc99 100644 --- a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -35,6 +35,9 @@ # 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 # Bound all serialized OffsetIndexes and their retained typed PageLocations. _MAX_INDEX_BYTES = 8 * 1024 * 1024 _MAX_PAGE_LOCATIONS = 128 * 1024 @@ -115,22 +118,22 @@ class _OffsetIndexDecoder: def __init__(self, data, max_locations): self.parser = _Compact(data) self.max_locations = max_locations - self.remaining_items = max_locations * 4 + # PageLocations retain three fields each. Allow the standard optional + # per-page byte-count list plus a small amount of forward metadata. + self.remaining_items = max_locations * 6 + 16 def decode(self): locations = None previous = 0 - seen = set() while True: field = self._field(previous) if field is None: break field_id, kind = field - if field_id in seen: - raise ValueError("Duplicate Parquet OffsetIndex field") - seen.add(field_id) 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() @@ -148,6 +151,7 @@ def _field(self, previous): 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): @@ -173,17 +177,15 @@ def _location(self): values = [None, None, None] expected = (6, 5, 6) previous = 0 - seen = set() while True: field = self._field(previous) if field is None: break field_id, kind = field - if field_id in seen: - raise ValueError("Duplicate Parquet PageLocation field") - seen.add(field_id) 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) @@ -245,15 +247,11 @@ def _skip(self, kind, depth=0): return if kind == 12: previous = 0 - seen = set() while True: field = self._field(previous) if field is None: return field_id, field_kind = field - if field_id in seen: - raise ValueError("Duplicate Parquet compact field") - seen.add(field_id) previous = field_id self._skip(field_kind, depth + 1) raise ValueError("Unsupported Parquet compact type: {}".format(kind)) @@ -365,10 +363,16 @@ def create(cls, source, parquet_file, columns, row_groups, batch_size): "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(" _MAX_FOOTER_BYTES: + return None footer = _Compact(serialized[-8 - length:-8]).value(12) if 8 in footer or 9 in footer: return None # Encrypted pages need the original file identity/AAD. diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 5b49e2c7a95c..91236c6459ee 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -114,7 +114,8 @@ def test_ranges_projection_missing_fields_and_fallback_in_same_file(fixture): @pytest.mark.parametrize( - 'mode', ['full', 'no_index', 'budget', 'location_budget', 'cache', 'scattered']) + 'mode', ['full', 'no_index', 'budget', 'location_budget', 'footer_bytes', + 'footer_chunks', 'cache', 'scattered']) def test_unsupported_or_expensive_reads_fall_back(fixture, mode): path, table, file_io, counter = fixture kwargs = {} @@ -129,6 +130,10 @@ def test_unsupported_or_expensive_reads_fall_back(fixture, mode): 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.ParquetPageIndexReader, '_batches', side_effect=AssertionError('must fall back')): result, _ = _read(fixture, **kwargs) @@ -216,6 +221,27 @@ def test_offset_index_page_locations_are_bounded_before_decoding(): 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_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)] From a77377aa4a2c0bc5c7d56d644c515e8022d448d7 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Sun, 20 Sep 2026 01:45:47 -0700 Subject: [PATCH 10/10] [python] Bound Parquet page header decoding --- .../read/reader/parquet_page_index_reader.py | 241 +++++++++++++----- .../pypaimon/tests/parquet_page_index_test.py | 40 ++- 2 files changed, 210 insertions(+), 71 deletions(-) diff --git a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py index bdf6283cbc99..6ca5ca25cacf 100644 --- a/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py +++ b/paimon-python/pypaimon/read/reader/parquet_page_index_reader.py @@ -38,9 +38,11 @@ # 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): @@ -50,9 +52,17 @@ class _PageIndexBudgetExceeded(Exception): class _Compact: """Thrift compact values used by Parquet metadata (no generated bindings).""" - def __init__(self, data): + 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 @@ -92,6 +102,7 @@ def value(self, kind, depth=0): 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) @@ -107,41 +118,18 @@ def value(self, kind, depth=0): 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 _OffsetIndexDecoder: - """Decode typed PageLocations without materializing generic Thrift trees.""" +class _BoundedCompactDecoder: + """Skip unknown compact fields without materializing object trees.""" - def __init__(self, data, max_locations): + def __init__(self, data, max_items): self.parser = _Compact(data) - self.max_locations = max_locations - # PageLocations retain three fields each. Allow the standard optional - # per-page byte-count list plus a small amount of forward metadata. - self.remaining_items = max_locations * 6 + 16 - - 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 + self.remaining_items = max_items def _field(self, previous): header = self.parser.take(1)[0] @@ -163,41 +151,6 @@ def _collection(self): raise ValueError("Invalid Parquet compact collection size") return count, element - 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 - def _consume(self, count): if count > self.remaining_items: raise _PageIndexBudgetExceeded( @@ -257,10 +210,138 @@ def _skip(self, kind, depth=0): 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: @@ -373,7 +454,11 @@ def create(cls, source, parquet_file, columns, row_groups, batch_size): length = struct.unpack(" _MAX_FOOTER_BYTES: return None - footer = _Compact(serialized[-8 - length:-8]).value(12) + 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] @@ -493,7 +578,23 @@ def read_row_group(self, group, runs): if (selected_bytes + index_bytes >= full_bytes or selected_bytes > _MAX_PAGE_BYTES): return None - return self._batches(plans, runs) + 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 @@ -512,9 +613,9 @@ def _column_payload(self, plan): num_values = 0 repeated = self.metadata.schema.column(index).max_repetition_level > 0 for position, (_, length) in enumerate(ranges): - parser = _Compact(memoryview(payload)[cursor:cursor + length]) - header = parser.value(12) - if parser.position + _get(header, 3) != length: + 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: @@ -536,7 +637,7 @@ def _column_payload(self, plan): if actual != expected or values < expected: raise ValueError("Parquet page rows disagree with OffsetIndex") num_values += values - uncompressed_size += parser.position + _get(header, 2) + uncompressed_size += header_size + _get(header, 2) cursor += length patched_column = {key: value for key, value in column.items() if key <= 8} diff --git a/paimon-python/pypaimon/tests/parquet_page_index_test.py b/paimon-python/pypaimon/tests/parquet_page_index_test.py index 91236c6459ee..8a33e0b1dc80 100644 --- a/paimon-python/pypaimon/tests/parquet_page_index_test.py +++ b/paimon-python/pypaimon/tests/parquet_page_index_test.py @@ -115,7 +115,7 @@ def test_ranges_projection_missing_fields_and_fallback_in_same_file(fixture): @pytest.mark.parametrize( 'mode', ['full', 'no_index', 'budget', 'location_budget', 'footer_bytes', - 'footer_chunks', 'cache', 'scattered']) + 'footer_chunks', 'footer_items', 'cache', 'scattered']) def test_unsupported_or_expensive_reads_fall_back(fixture, mode): path, table, file_io, counter = fixture kwargs = {} @@ -134,6 +134,8 @@ def test_unsupported_or_expensive_reads_fall_back(fixture, mode): 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) @@ -231,6 +233,42 @@ def test_offset_index_unknown_struct_fields_are_bounded(): 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: