Skip to content
Draft
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
7 changes: 7 additions & 0 deletions docs/docs/pypaimon/blob.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,13 @@ Without `blob-as-descriptor=true`, blob values are materialized before
`row.get_blob(...)` returns; `new_input_stream()` then reads from
in-memory bytes, not from storage.

For data-evolution reads, PyPaimon applies user filters, row-level authorization
filters, and limits before materializing projected scalar BLOB payloads. A user
or authorization filter that references a BLOB value keeps that field eager.
Column masking is applied after payload materialization. Set
`read.defer-blob-resolve=false` to restore eager materialization. ARRAY and MAP
elements containing BLOB values are not deferred.

## Lower-level: `Blob.from_bytes`

When you already have raw or descriptor bytes (for example from a custom
Expand Down
4 changes: 4 additions & 0 deletions docs/docs/pypaimon/pytorch.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,10 @@ when it is false, it will read the full amount of data into memory.

**`prefetch_concurrency`** (default: 1): When streaming is true, number of threads used for parallel prefetch within each DataLoader worker. Set to a value greater than 1 to partition splits across threads and increase read throughput. Has no effect when streaming is false.

When the read builder has a `LIMIT`, streaming reads use one DataLoader worker
and disable split prefetch fan-out. This preserves one global remaining-row
quota and avoids materializing BLOB payloads beyond the limit.

## Shuffle

PyPaimon supports streaming shuffle for PyTorch `IterableDataset`. The shuffle
Expand Down
13 changes: 13 additions & 0 deletions paimon-python/pypaimon/common/options/core_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -903,6 +903,16 @@ class CoreOptions:
)
)

READ_DEFER_BLOB_RESOLVE: ConfigOption[bool] = (
ConfigOptions.key("read.defer-blob-resolve")
.boolean_type()
.default_value(True)
.with_description(
"Whether filtered or limited data-evolution reads should apply "
"row selection before materializing projected scalar BLOB payloads."
)
)

READ_BATCH_SIZE: ConfigOption[int] = (
ConfigOptions.key("read.batch-size")
.int_type()
Expand Down Expand Up @@ -1526,6 +1536,9 @@ def local_cache_block_size(self) -> MemorySize:
def local_cache_whitelist(self) -> str:
return self.options.get(CoreOptions.LOCAL_CACHE_WHITELIST)

def read_defer_blob_resolve(self) -> bool:
return self.options.get(CoreOptions.READ_DEFER_BLOB_RESOLVE)

def read_batch_size(self, default=None) -> int:
return self.options.get(CoreOptions.READ_BATCH_SIZE, default or 1024)

Expand Down
10 changes: 8 additions & 2 deletions paimon-python/pypaimon/read/datasource/torch_dataset.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,11 @@ def _row_to_dict(self, offset_row) -> dict:
return row_dict

def _worker_splits(self, worker_info) -> List[Split]:
if self.table_read.limit is not None:
if worker_info is None or worker_info.id == 0:
return self.splits
return []

if worker_info is None:
return self.splits

Expand Down Expand Up @@ -164,7 +169,7 @@ def __iter__(self):
worker_info = torch.utils.data.get_worker_info()
splits_to_process = self._worker_splits(worker_info)

if self.prefetch_concurrency > 1:
if self.prefetch_concurrency > 1 and self.table_read.limit is None:
for row in self._iter_rows(splits_to_process):
yield row
return
Expand Down Expand Up @@ -288,7 +293,8 @@ def __iter__(self):
worker_id = worker_info.id if worker_info is not None else 0
splits_to_process = self._worker_splits(worker_info)

if self.max_buffer_input_splits == 1:
if (self.table_read.limit is not None
or self.max_buffer_input_splits == 1):
rows = self._iter_ordered_rows(splits_to_process)
else:
rows = self._iter_interleaved_rows(splits_to_process)
Expand Down
75 changes: 75 additions & 0 deletions paimon-python/pypaimon/read/reader/deferred_blob_resolve_reader.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
# 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.

from typing import List, Optional

import pyarrow as pa
from pyarrow import RecordBatch

from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader
from pypaimon.table.row.blob import Blob


class DeferredBlobResolveReader(RecordBatchReader):
"""Materialize projected BLOB payloads after row filtering.

This must remain the outermost BLOB materialization layer because adopted
metadata still identifies the materialized columns as logical BLOB fields.
"""

def __init__(self, inner: RecordBatchReader, file_io,
blob_field_names: List[str], blob_parallelism: int = 1):
self._inner = inner
self._file_io = file_io
self._blob_field_names = blob_field_names
self._blob_parallelism = max(1, blob_parallelism)
self._adopt_metadata(inner)

def read_arrow_batch(self) -> Optional[RecordBatch]:
batch = self._inner.read_arrow_batch()
if batch is None:
return None

columns = list(batch.columns)
fields = list(batch.schema)
changed = False
for field_name in self._blob_field_names:
column_index = batch.schema.get_field_index(field_name)
if column_index < 0:
continue
values = batch.column(column_index).to_pylist()
blobs = [Blob.from_bytes(value, self._file_io) for value in values]
payloads = self._file_io.read_blobs_concurrent(
blobs, self._blob_parallelism)
source_field = batch.schema.field(column_index)
columns[column_index] = pa.array(payloads, type=pa.large_binary())
fields[column_index] = pa.field(
field_name,
pa.large_binary(),
nullable=source_field.nullable,
metadata=source_field.metadata,
)
changed = True
if not changed:
return batch
return pa.RecordBatch.from_arrays(
columns,
schema=pa.schema(fields, metadata=batch.schema.metadata),
)

def close(self) -> None:
self._inner.close()
93 changes: 83 additions & 10 deletions paimon-python/pypaimon/read/split_read.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,10 @@
MergeAllBatchReader, DataEvolutionMergeReader)
from pypaimon.read.reader.concat_record_reader import ConcatRecordReader

from pypaimon.read.reader.auth_masking_reader import AuthFilterReader
from pypaimon.read.reader.data_file_batch_reader import DataFileBatchReader
from pypaimon.read.reader.deferred_blob_resolve_reader import \
DeferredBlobResolveReader
from pypaimon.read.reader.drop_delete_reader import DropDeleteRecordReader
from pypaimon.read.reader.empty_record_reader import EmptyFileRecordReader
from pypaimon.read.reader.field_bunch import BlobBunch, DataBunch, FieldBunch, VectorBunch
Expand Down Expand Up @@ -81,6 +84,33 @@
KEY_PREFIX = "_KEY_"
KEY_FIELD_ID_START = 1000000
NULL_FIELD_INDEX = -1


def deferred_blob_field_names(table, read_fields: List[DataField],
predicate: Optional[Predicate],
limit: Optional[int],
has_post_filter: bool = False) -> set:
# An auth filter also selects rows; defer past it too, like a predicate/limit.
if ((predicate is None and limit is None and not has_post_filter)
or CoreOptions.blob_as_descriptor(table.options)
or not table.options.read_defer_blob_resolve()):
return set()

inline_fields = (
CoreOptions.blob_descriptor_fields(table.options)
| CoreOptions.blob_view_fields(table.options)
)
predicate_fields = (
predicate_field_names(predicate) if predicate is not None else set()
)
return {
read_fields[index].name
for index in blob_field_indices(read_fields)
if read_fields[index].name not in inline_fields
and read_fields[index].name not in predicate_fields
}


ROW_SIDECAR_FORMAT = CoreOptions.FILE_FORMAT_ROW

_COMPRESS_EXTENSIONS = frozenset(['gz', 'bz2', 'deflate', 'snappy', 'lz4', 'zst'])
Expand Down Expand Up @@ -283,7 +313,7 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool,
if has_nested:
raise NotImplementedError(
"Nested-field projection is not supported on BLOB files")
blob_as_descriptor = CoreOptions.blob_as_descriptor(self.table.options)
blob_as_descriptor = self._read_blob_as_descriptor(read_file_fields)
blob_parallelism = self._blob_parallelism
format_reader = FormatBlobReader(self.table.file_io, file_path, read_file_fields,
self.read_fields, read_arrow_predicate, blob_as_descriptor,
Expand Down Expand Up @@ -419,6 +449,12 @@ def file_reader_supplier(self, file: DataFileMeta, for_merge_read: bool,

return reader

def _read_blob_as_descriptor(self, field_names: List[str]) -> bool:
if CoreOptions.blob_as_descriptor(self.table.options):
return True
deferred_fields = getattr(self, '_deferred_blob_fields', set())
return any(field_name in deferred_fields for field_name in field_names)

@staticmethod
def _row_sidecar_file_name(file: DataFileMeta) -> Optional[str]:
row_files = [
Expand Down Expand Up @@ -995,7 +1031,10 @@ def __init__(
nested_name_paths: Optional[List[List[str]]] = None,
limit: Optional[int] = None,
outer_extract_name_paths: Optional[List[List[str]]] = None,
outer_flat_read_type: Optional[List[DataField]] = None):
outer_flat_read_type: Optional[List[DataField]] = None,
post_merge_filter=None,
eager_blob_fields=None,
post_filter_after_inline=False):
self.row_ranges = None
actual_split = split
if isinstance(split, IndexedSplit):
Expand All @@ -1008,6 +1047,11 @@ def __init__(
)
self.outer_extract_name_paths = outer_extract_name_paths
self.outer_flat_read_type = outer_flat_read_type
self._post_merge_filter = post_merge_filter
# Apply the auth filter after inline BLOB resolution, so scalar BLOBs still defer.
self._post_filter_after_inline = post_filter_after_inline
self._eager_blob_fields = set(eager_blob_fields or [])
self._deferred_blob_fields = self._deferred_blob_field_names()

def _push_down_predicate(self) -> Optional[Predicate]:
# Data evolution: files may have different schemas, so we don't push predicate
Expand All @@ -1027,8 +1071,35 @@ def create_reader(self) -> RecordReader:
prescan_reader_factory=lambda names: self._create_prescan_reader(names),
blob_parallelism=blob_parallelism)

if self._post_filter_after_inline:
if self._post_merge_filter is not None:
reader = AuthFilterReader(reader, self._post_merge_filter)
if self.limit is not None:
reader = LimitedRecordBatchReader(reader, self.limit)

if self._deferred_blob_fields:
blob_names = [
field.name for field in self.read_fields
if field.name in self._deferred_blob_fields
]
reader = DeferredBlobResolveReader(
reader,
self.table.file_io,
blob_names,
blob_parallelism=self._blob_parallelism,
)

return reader

def _deferred_blob_field_names(self) -> set:
return deferred_blob_field_names(
self.table,
self.read_fields,
self.predicate_for_reader,
self.limit,
has_post_filter=self._post_merge_filter is not None,
) - self._eager_blob_fields

def _create_raw_reader(self) -> RecordReader:
"""Core read logic: split_by_row_id -> suppliers -> ConcatBatchReader -> filter."""
files = self.split.files
Expand Down Expand Up @@ -1065,6 +1136,9 @@ def _create_raw_reader(self) -> RecordReader:
else:
reader = merge_reader

if self._post_merge_filter is not None and not self._post_filter_after_inline:
reader = AuthFilterReader(reader, self._post_merge_filter)

if self.outer_extract_name_paths:
if self.outer_flat_read_type is None:
raise ValueError(
Expand All @@ -1075,7 +1149,7 @@ def _create_raw_reader(self) -> RecordReader:
reader = NestedLeafBatchReader(
reader, self.outer_extract_name_paths, self.outer_flat_read_type)

if self.limit is not None:
if self.limit is not None and not self._post_filter_after_inline:
reader = LimitedRecordBatchReader(reader, self.limit)

return reader
Expand Down Expand Up @@ -1138,17 +1212,16 @@ def _create_prescan_reader(self, field_names):
if not prescan_fields:
return EmptyRecordBatchReader()

# When there's a normal field predicate, don't push down limit to prescan reader
# because the outer reader will apply predicate+limit filtering,
# while prescan reader would only apply limit without normal field predicate
# TODO support limit+predicate push down
# Skip limit push-down when the outer reader also selects rows (predicate or auth
# filter): prescan's first-N rows would differ from the outer set. TODO: push down.
skip_limit = self.predicate is not None or self._post_merge_filter is not None
prescan_read = DataEvolutionSplitRead(
table=self.table,
predicate=self.predicate,
read_type=prescan_fields,
split=self.split,
row_tracking_enabled=False,
limit=None if self.predicate else self.limit,
limit=None if skip_limit else self.limit,
)
prescan_read.row_ranges = self.row_ranges
return prescan_read._create_raw_reader()
Expand Down Expand Up @@ -1297,7 +1370,7 @@ def _create_union_reader(self, need_merge_files: List[DataFileMeta], deletion_ve
[read_fields[0]]
).field(0).type,
self.row_ranges,
CoreOptions.blob_as_descriptor(self.table.options),
self._read_blob_as_descriptor([read_fields[0].name]),
deletion_vector=deletion_vector,
batch_size=batch_size,
blob_parallelism=self._blob_parallelism,
Expand Down Expand Up @@ -1352,7 +1425,7 @@ def _create_raw_blob_file_reader(
read_fields,
self.read_fields,
None,
CoreOptions.blob_as_descriptor(self.table.options),
self._read_blob_as_descriptor(read_fields),
batch_size=self.table.options.read_batch_size(),
row_indices=row_indices,
blob_parallelism=blob_parallelism,
Expand Down
Loading
Loading