Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions paimon-python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,19 @@ pip3 install dist/*.tar.gz

The command will install the package and core dependencies to your local Python environment.

# Parquet page-index reads

For row-tracking tables with a Parquet OffsetIndex, PyPaimon can read a
contiguous `_ROW_ID` range without decoding the full row group. This is enabled
by default and can be disabled with the table option:

```python
table = table.copy({"parquet.filter.columnindex.enabled": "false"})
```

Unsupported reads use the normal path. Reading fewer bytes may require more
object-store requests.

# Native scan planning

PyPaimon can plan splits with the optional `pypaimon-rust` package while retaining
Expand Down
16 changes: 16 additions & 0 deletions paimon-python/pypaimon/common/options/core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -1128,6 +1128,19 @@ class CoreOptions:
.with_description("Read batch size for any file format if it supports.")
)

PARQUET_COLUMN_INDEX_ENABLED: ConfigOption[bool] = (
ConfigOptions.key("parquet.filter.columnindex.enabled")
.boolean_type()
.default_value(True)
.with_description(
"Enable Parquet page-index pruning. PyPaimon currently uses OffsetIndex "
"metadata for contiguous row windows. "
"Requires existing offset indexes; nested fields use common leaf row boundaries. "
"Unsupported or expensive selections use the ordinary reader. "
"Does not enable ColumnIndex predicate filtering."
)
)

READ_PARALLELISM: ConfigOption[int] = (
ConfigOptions.key("read.parallelism")
.int_type()
Expand Down Expand Up @@ -1878,6 +1891,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 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)

Expand Down
66 changes: 64 additions & 2 deletions paimon-python/pypaimon/read/reader/format_pyarrow_reader.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -526,15 +528,37 @@ 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.parquet_column_index_enabled()
and self._row_group_cache is None
and not self._bounded_variant_read):
from pypaimon.read.reader.parquet_page_index_reader import (
ParquetPageIndexReader,
)
self._page_index_reader = ParquetPageIndexReader.create(
self._parquet_source, self._parquet_file,
self._row_group_read_columns(),
self._selected_parquet_row_groups, batch_size)
except BaseException:
self._parquet_source.close()
raise
if file_format == 'orc' and self._selected_shared_map_paths:
import pyarrow.orc as orc
self._orc_source = file_io.filesystem.open_input_file(
file_path_for_pyarrow)
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:
Expand Down Expand Up @@ -611,6 +635,37 @@ def _iter_row_group_batches(self):
if out.num_rows:
yield out

def _iter_page_index_batches(self, selected_infos, runs):
select = self._select_nested_fields if self._has_nested_path else self._select_existing_fields
run_index = 0
for group, (offset, count) in zip(
self._selected_parquet_row_groups, selected_infos):
while run_index < len(runs) and runs[run_index][1] < offset:
run_index += 1
local_runs = []
position = run_index
while position < len(runs) and runs[position][0] < offset + count:
lower, upper = runs[position]
local_runs.append((max(0, lower - offset),
min(count - 1, upper - offset)))
position += 1
batches = self._page_index_reader.read_row_group(group, local_runs)
if batches is None:
raw = self._read_parquet_row_group_batches(
group, self._row_group_read_columns())
slicer = _RowRunSlicer([(0, count)], local_runs)
while True:
batch = slicer.next_batch(raw)
if batch is None:
break
yield select(batch)
else:
try:
for batch in batches:
yield select(batch)
finally:
batches.close()

def _read_parquet_row_group_batches(self, row_group, columns):
return self._parquet_file.iter_batches(
row_groups=[row_group],
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading