From 743fd8c47b1e2d76c9cd428041bd03069562b5fa Mon Sep 17 00:00:00 2001 From: "wenchao.wu" Date: Thu, 27 Aug 2026 16:22:02 +0800 Subject: [PATCH] [python] Route descriptor-backed BLOB reads through the table FileIO. Reuse UriReaderFactory.from_file_io for non-HTTP URIs, resolve known descriptor fields with from_descriptor_bytes (including v1 writes), and keep RowKind/LIMIT stable around inline convert. Do not pin FileIO in the from_file_io cache. --- .../catalog/rest/rest_token_file_io.py | 5 +- paimon-python/pypaimon/common/file_io.py | 57 +- paimon-python/pypaimon/common/uri_reader.py | 119 +- .../filesystem/hdfs_native_file_io.py | 3 + .../pypaimon/filesystem/local_file_io.py | 5 + .../pypaimon/filesystem/pyarrow_file_io.py | 5 + .../read/reader/auth_masking_reader.py | 40 +- .../reader/blob_descriptor_convert_reader.py | 95 +- .../read/reader/blob_view_read_support.py | 68 + .../read/reader/concat_batch_reader.py | 6 +- .../reader/deferred_blob_resolve_reader.py | 5 + .../pypaimon/read/reader/field_indices.py | 22 + .../read/reader/filter_record_batch_reader.py | 3 + .../read/reader/filter_record_reader.py | 2 + .../read/reader/iface/record_batch_reader.py | 17 +- .../read/reader/iface/record_reader.py | 15 + .../read/reader/limited_record_reader.py | 3 + .../read/reader/nested_leaf_batch_reader.py | 8 +- .../reader/outer_projection_record_reader.py | 19 +- .../reader/row_range_filter_record_reader.py | 1 + paimon-python/pypaimon/read/split_read.py | 173 ++- paimon-python/pypaimon/read/table_read.py | 9 + .../pypaimon/table/row/offset_row.py | 56 +- .../pypaimon/tests/blob_table_test.py | 187 +++ paimon-python/pypaimon/tests/blob_test.py | 1192 +++++++++++++++++ .../pypaimon/tests/resolving_file_io_test.py | 34 + .../tests/rest/rest_token_file_io_test.py | 20 + .../pypaimon/tests/uri_reader_factory_test.py | 91 ++ .../pypaimon/tests/vector_table_test.py | 7 + .../pypaimon/utils/blob_view_lookup.py | 41 +- 30 files changed, 2187 insertions(+), 121 deletions(-) create mode 100644 paimon-python/pypaimon/read/reader/blob_view_read_support.py diff --git a/paimon-python/pypaimon/catalog/rest/rest_token_file_io.py b/paimon-python/pypaimon/catalog/rest/rest_token_file_io.py index 42dabb268cb2..5e51861a1efe 100644 --- a/paimon-python/pypaimon/catalog/rest/rest_token_file_io.py +++ b/paimon-python/pypaimon/catalog/rest/rest_token_file_io.py @@ -299,4 +299,7 @@ def valid_token(self): return self.token def close(self): - pass + factory = self._uri_reader_factory_cache + self._uri_reader_factory_cache = None + if factory is not None: + factory.close() diff --git a/paimon-python/pypaimon/common/file_io.py b/paimon-python/pypaimon/common/file_io.py index 10406849b341..168a84bd2fd2 100644 --- a/paimon-python/pypaimon/common/file_io.py +++ b/paimon-python/pypaimon/common/file_io.py @@ -418,26 +418,51 @@ def _run_lane(lane): def read_blobs_concurrent(self, blobs, parallelism): """Read a list of Blobs concurrently, coalescing same-file ranged reads. - ``BlobRef`` values expose a file range and are coalesced; in-memory - ``BlobData`` values are returned directly. + Exact ``BlobRef`` values (not subclasses) with a file-backed UriReader + are coalesced through that FileIO so table-scoped credentials are + preserved. Subclasses may override ``new_input_stream()`` and must not + be bypassed. Other readers (for example HTTP) read through the Blob. """ - from pypaimon.table.row.blob import BlobRef + from concurrent.futures import ThreadPoolExecutor + + from pypaimon.common.uri_reader import FileUriReader + from pypaimon.table.row.blob import BlobData, BlobRef + results: List[Optional[bytes]] = [None] * len(blobs) - ranges: List[Optional[tuple]] = [None] * len(blobs) - inmem = [] - for i, b in enumerate(blobs): - if b is None: + file_groups = {} + other_blobs = [] + for index, blob in enumerate(blobs): + if blob is None: continue - if isinstance(b, BlobRef): - d = b.to_descriptor() - ranges[i] = (d.uri, d.offset, d.length) + if isinstance(blob, BlobData): + results[index] = blob.to_data() + elif type(blob) is BlobRef and isinstance( + blob.uri_reader, FileUriReader): + descriptor = blob.to_descriptor() + source_file_io = blob.uri_reader.file_io + group = file_groups.setdefault( + id(source_file_io), (source_file_io, []))[1] + group.append((index, ( + descriptor.uri, descriptor.offset, descriptor.length))) else: - inmem.append((i, b)) - for i, v in enumerate(self.read_ranges_coalesced(ranges, parallelism)): - if v is not None: - results[i] = v - for idx, b in inmem: - results[idx] = b.to_data() + other_blobs.append((index, blob)) + + for source_file_io, indexed_ranges in file_groups.values(): + ranges = [value for _, value in indexed_ranges] + values = source_file_io.read_ranges_coalesced(ranges, parallelism) + for (index, _), value in zip(indexed_ranges, values): + results[index] = value + + if other_blobs: + workers = max(1, min(parallelism, len(other_blobs))) + + def _read_blob(indexed_blob): + return indexed_blob[1].to_data() + + with ThreadPoolExecutor(workers) as pool: + values = pool.map(_read_blob, other_blobs) + for (index, _), value in zip(other_blobs, values): + results[index] = value return results def read_file_utf8(self, path: str) -> str: diff --git a/paimon-python/pypaimon/common/uri_reader.py b/paimon-python/pypaimon/common/uri_reader.py index a417a0c9ae41..5548329c2965 100644 --- a/paimon-python/pypaimon/common/uri_reader.py +++ b/paimon-python/pypaimon/common/uri_reader.py @@ -16,6 +16,7 @@ # under the License. import io +import weakref from abc import ABC, abstractmethod from typing import Any, Optional, Union from urllib.parse import urlparse, ParseResult @@ -57,6 +58,10 @@ class FileUriReader(UriReader): def __init__(self, file_io: Any): self._file_io = file_io + @property + def file_io(self) -> Any: + return self._file_io + def new_input_stream(self, uri: str): try: return self._file_io.new_input_stream(uri) @@ -109,8 +114,33 @@ class UriReaderFactory: def __init__(self, catalog_options: Union[Options, dict]) -> None: self.catalog_options = catalog_options if isinstance(catalog_options, Options) else Options(catalog_options) - self._readers = LRUCache(CatalogOptions.BLOB_FILE_IO_DEFAULT_CACHE_SIZE) self._readers_lock = rwlock.RWLockFair() + # FileIOs created by this factory. Do not close them on LRU eviction: + # live BlobRefs may still hold the corresponding UriReader. + self._owned_file_ios = [] + self._closing = False + self._readers = self._new_reader_cache() + + _FROM_FILE_IO_FACTORIES = weakref.WeakKeyDictionary() + + @staticmethod + def from_file_io(file_io: Any) -> 'UriReaderFactory': + """Reuse a token-aware FileIO for non-HTTP URIs (Java fromFileIO).""" + try: + cached = UriReaderFactory._FROM_FILE_IO_FACTORIES.get(file_io) + except TypeError: + return _ProvidedFileIOUriReaderFactory(file_io) + if cached is not None: + return cached + factory = _ProvidedFileIOUriReaderFactory(file_io) + try: + UriReaderFactory._FROM_FILE_IO_FACTORIES[file_io] = factory + except TypeError: + pass + return factory + + def _new_reader_cache(self) -> LRUCache: + return LRUCache(CatalogOptions.BLOB_FILE_IO_DEFAULT_CACHE_SIZE) def create(self, input_uri: str) -> UriReader: try: @@ -148,12 +178,38 @@ def _new_reader(self, key: UriKey, parsed_uri: ParseResult) -> UriReader: from pypaimon.common.file_io import FileIO uri_string = parsed_uri.geturl() file_io = FileIO.get(uri_string, self.catalog_options) + self._owned_file_ios.append(file_io) return UriReader.from_file(file_io) except Exception as e: raise RuntimeError(f"Failed to create reader for URI {parsed_uri.geturl()}") from e def clear_cache(self) -> None: - self._readers.clear() + if self._closing: + return + self._closing = True + wlock = self._readers_lock.gen_wlock() + wlock.acquire() + try: + file_ios = list(self._owned_file_ios) + self._owned_file_ios = [] + self._readers = self._new_reader_cache() + finally: + wlock.release() + first_error = None + try: + for file_io in file_ios: + try: + file_io.close() + except Exception as error: + if first_error is None: + first_error = error + finally: + self._closing = False + if first_error is not None: + raise first_error + + def close(self) -> None: + self.clear_cache() def get_cache_size(self) -> int: return len(self._readers) @@ -161,8 +217,67 @@ def get_cache_size(self) -> int: def __getstate__(self): state = self.__dict__.copy() del state['_readers_lock'] + del state['_readers'] + del state['_owned_file_ios'] return state def __setstate__(self, state): self.__dict__.update(state) self._readers_lock = rwlock.RWLockFair() + self._owned_file_ios = [] + self._closing = False + self._readers = self._new_reader_cache() + + +class _ProvidedFileIOUriReaderFactory(UriReaderFactory): + """Resolves HTTP(S) via HttpUriReader and every other URI through file_io.""" + + def __init__(self, file_io: Any) -> None: + super().__init__({}) + self._bind_provided_file_io(file_io) + + def _bind_provided_file_io(self, file_io: Any) -> None: + try: + self._provided_file_io = weakref.ref(file_io) + except TypeError: + # Not weakref-able, and therefore also not a WeakKeyDictionary + # key — from_file_io does not cache these objects. + self._provided_file_io = lambda: file_io + + def _resolved_file_io(self): + file_io = self._provided_file_io() + if file_io is None: + raise RuntimeError( + "FileIO used by UriReaderFactory.from_file_io was garbage collected") + return file_io + + def __getstate__(self): + state = super().__getstate__() + # weakref.ref (and the TypeError fallback lambda) cannot be pickled. + # Resolve to a strong FileIO for the wire; __setstate__ re-wraps. + state['_provided_file_io'] = self._resolved_file_io() + return state + + def __setstate__(self, state): + file_io = state.pop('_provided_file_io') + super().__setstate__(state) + self._bind_provided_file_io(file_io) + + def create(self, input_uri: str) -> UriReader: + try: + parsed_uri = urlparse(input_uri) + except Exception as e: + raise ValueError("Invalid URI: %s" % input_uri) from e + scheme = (parsed_uri.scheme or '').lower() + if scheme in ('http', 'https'): + return super().create(input_uri) + # Do not LRU-cache FileUriReader: it holds FileIO strongly and would + # pin the WeakKeyDictionary key. Every non-HTTP URI already wraps the + # same provided FileIO, so the cache buys nothing here. + return UriReader.from_file(self._resolved_file_io()) + + def _new_reader(self, key: UriKey, parsed_uri: ParseResult) -> UriReader: + scheme = (key.scheme or '').lower() + if scheme in ('http', 'https'): + return UriReader.from_http() + return UriReader.from_file(self._resolved_file_io()) diff --git a/paimon-python/pypaimon/filesystem/hdfs_native_file_io.py b/paimon-python/pypaimon/filesystem/hdfs_native_file_io.py index 78cad2140e28..627ea7afbdc3 100644 --- a/paimon-python/pypaimon/filesystem/hdfs_native_file_io.py +++ b/paimon-python/pypaimon/filesystem/hdfs_native_file_io.py @@ -696,4 +696,7 @@ def write_vortex(self, path: str, data: pyarrow.Table, **kwargs): raise RuntimeError(f"Failed to write Vortex file {path}: {e}") from e def close(self): + uri_reader_factory = getattr(self, 'uri_reader_factory', None) + if uri_reader_factory is not None: + uri_reader_factory.close() self._client = None diff --git a/paimon-python/pypaimon/filesystem/local_file_io.py b/paimon-python/pypaimon/filesystem/local_file_io.py index e3530c1951c2..ec70e3f7ec4b 100644 --- a/paimon-python/pypaimon/filesystem/local_file_io.py +++ b/paimon-python/pypaimon/filesystem/local_file_io.py @@ -490,6 +490,11 @@ def write_blob(self, path: str, data: pyarrow.Table, **kwargs): self.delete_quietly(path) raise RuntimeError(f"Failed to write blob file {path}: {e}") from e + def close(self): + uri_reader_factory = getattr(self, 'uri_reader_factory', None) + if uri_reader_factory is not None: + uri_reader_factory.close() + class FuseLocalFileIO(LocalFileIO): """LocalFileIO that translates remote OSS paths to FUSE-mounted local paths. diff --git a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py index 362994fb65ef..1cf4540e48ee 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -104,6 +104,11 @@ def __setstate__(self, state): self.__dict__.update(state) self._legacy_bucket_lock = threading.Lock() + def close(self): + uri_reader_factory = getattr(self, 'uri_reader_factory', None) + if uri_reader_factory is not None: + uri_reader_factory.close() + @staticmethod def parse_location(location: str): uri = urlparse(location) diff --git a/paimon-python/pypaimon/read/reader/auth_masking_reader.py b/paimon-python/pypaimon/read/reader/auth_masking_reader.py index aeee5e13e92f..a05024216ec3 100644 --- a/paimon-python/pypaimon/read/reader/auth_masking_reader.py +++ b/paimon-python/pypaimon/read/reader/auth_masking_reader.py @@ -40,7 +40,10 @@ def __init__(self, inner, schema: pa.Schema, chunk_size: int = 65536, include_ro self._exhausted = False self._pending_iterator = None self._include_row_kind = include_row_kind + self.file_io = getattr(inner, 'file_io', None) self.blob_field_indices = getattr(inner, 'blob_field_indices', None) + self.descriptor_field_indices = getattr(inner, 'descriptor_field_indices', None) + self.blob_view_lookup = getattr(inner, 'blob_view_lookup', None) self.vector_field_indices = getattr(inner, 'vector_field_indices', None) def read_arrow_batch(self) -> Optional[pa.RecordBatch]: @@ -66,6 +69,7 @@ def read_arrow_batch(self) -> Optional[pa.RecordBatch]: self._exhausted = True break self._pending_iterator = row_iterator + self._refresh_blob_view_lookup(self._inner) if not row_tuples: return None @@ -95,12 +99,25 @@ class BatchToRecordReaderAdapter(RecordReader): def __init__(self, inner: RecordBatchReader): self._inner = inner + self.file_io = getattr(inner, 'file_io', None) + self.blob_field_indices = getattr(inner, 'blob_field_indices', None) + self.descriptor_field_indices = getattr(inner, 'descriptor_field_indices', None) + self.blob_view_lookup = getattr(inner, 'blob_view_lookup', None) + self.vector_field_indices = getattr(inner, 'vector_field_indices', None) def read_batch(self): batch = self._inner.read_arrow_batch() if batch is None: return None - return _ArrowBatchIterator(batch) + self._refresh_blob_view_lookup(self._inner) + return _ArrowBatchIterator( + batch, + file_io=self.file_io, + blob_field_indices=self.blob_field_indices, + descriptor_field_indices=self.descriptor_field_indices, + blob_view_lookup=self.blob_view_lookup, + vector_field_indices=self.vector_field_indices, + ) def close(self): self._inner.close() @@ -108,7 +125,10 @@ def close(self): class _ArrowBatchIterator(RecordIterator): - def __init__(self, batch: pa.RecordBatch): + def __init__(self, batch: pa.RecordBatch, + file_io=None, blob_field_indices=None, + descriptor_field_indices=None, blob_view_lookup=None, + vector_field_indices=None): self._batch = batch self._idx = 0 self._has_rk = "_row_kind" in batch.schema.names @@ -118,6 +138,11 @@ def __init__(self, batch: pa.RecordBatch): else: self._rk_idx = -1 self._data_cols = list(range(batch.num_columns)) + self._file_io = file_io + self._blob_field_indices = blob_field_indices + self._descriptor_field_indices = descriptor_field_indices + self._blob_view_lookup = blob_view_lookup + self._vector_field_indices = vector_field_indices def next(self): if self._idx >= self._batch.num_rows: @@ -126,7 +151,13 @@ def next(self): self._batch.column(j)[self._idx].as_py() for j in self._data_cols ) - row = OffsetRow(row_tuple, 0, len(self._data_cols)) + row = OffsetRow( + row_tuple, 0, len(self._data_cols), + file_io=self._file_io, + blob_field_indices=self._blob_field_indices, + descriptor_field_indices=self._descriptor_field_indices, + blob_view_lookup=self._blob_view_lookup, + vector_field_indices=self._vector_field_indices) if self._has_rk: from pypaimon.table.row.row_kind import RowKind kind_str = self._batch.column(self._rk_idx)[self._idx].as_py() @@ -146,6 +177,7 @@ def read_arrow_batch(self) -> Optional[pa.RecordBatch]: batch = self._inner.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) mask = self._filter_fn(batch) return batch.filter(mask) @@ -184,6 +216,7 @@ def read_arrow_batch(self) -> Optional[pa.RecordBatch]: batch = self._inner.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) original_batch = batch masked_columns = {} for col_name, transform in self._parsed_rules.items(): @@ -223,6 +256,7 @@ def read_arrow_batch(self) -> Optional[pa.RecordBatch]: batch = self._inner.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) columns = self._columns if "_row_kind" in batch.schema.names and "_row_kind" not in columns: columns = ["_row_kind"] + list(columns) diff --git a/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py b/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py index 48cee2798850..f7bc46612f93 100644 --- a/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py +++ b/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py @@ -31,10 +31,9 @@ class BlobInlineConvertReader(RecordBatchReader): Processing is split into two clear stages: Stage 1 (BlobView resolution): If view fields exist, use a lightweight prescan reader (only projecting view columns) to collect - BlobViewStructs, bulk-preload their descriptors, then read - full data from the main reader and replace view field values - with descriptor bytes or real blob data according to the - blob-as-descriptor option. + BlobViewStructs and bulk-preload their descriptors, then replace + view field values with descriptor bytes so Stage 2 can + materialize payloads with the originating table FileIO. Stage 2 (BlobDescriptor resolution): Controlled by blob-as-descriptor option. If false, resolve BlobDescriptor bytes from descriptor fields into real blob data bytes. BlobView fields are already resolved @@ -67,24 +66,27 @@ def __init__(self, inner: RecordBatchReader, table, self._view_fields = CoreOptions.blob_view_fields(table.options) if resolve_enabled else set() self._descriptor_fields = CoreOptions.blob_descriptor_fields(table.options) self._blob_as_descriptor = CoreOptions.blob_as_descriptor(table.options) + if not self._blob_as_descriptor: + # Stage 2 materializes descriptor/view fields to payload bytes. + # Row-level descriptor routing must not re-parse that content. + self.descriptor_field_indices = set() self._prescan_done = False self._blob_view_lookup = None def read_arrow_batch(self) -> Optional[RecordBatch]: - # Align with Java: only enter blob view resolution when catalog_loader is available - # If catalog_loader is None, skip both Stage 1 (view resolution) and Stage 2 (descriptor resolution) + # Align with Java: only enter blob view resolution when catalog_loader is available. if self._view_fields and not self._prescan_done: self._prescan_view_structs() batch = self._inner.read_arrow_batch() if batch is None: return None - # Resolve view fields using the preloaded lookup - view_file_ios = {} + # Resolve view fields using the preloaded lookup. + view_blobs = {} if self._view_fields and self._blob_view_lookup is not None: - batch, view_file_ios = self._resolve_view_fields(batch, self._blob_view_lookup) + batch, view_blobs = self._resolve_view_fields(batch, self._blob_view_lookup) # Resolve BlobDescriptor -> real bytes (if blob-as-descriptor=false) - return self._resolve_descriptor_fields(batch, view_file_ios) + return self._resolve_descriptor_fields(batch, view_blobs) # ------------------------------------------------------------------ # Stage 1: BlobView prescan (lightweight, only reads view columns) @@ -125,33 +127,35 @@ def _prescan_view_structs(self): if all_view_structs: self._blob_view_lookup = BlobViewLookup(self._table) self._blob_view_lookup.preload(all_view_structs) + # Expose after prescan so OffsetRow.get_blob() can resolve each + # BlobViewStruct with the originating table FileIO. + self.blob_view_lookup = self._blob_view_lookup self._prescan_done = True def _resolve_view_fields(self, batch, blob_view_lookup): """Replace BlobViewStruct bytes in view fields with descriptor bytes.""" - view_file_ios = {} + view_blobs = {} for field_name in self._view_fields: if field_name not in batch.schema.names: continue values = [self._normalize_blob_to_bytes(v) for v in batch.column(field_name).to_pylist()] converted_values = [] - field_file_ios = [] + field_blobs = [] for value in values: if value is None or not ( isinstance(value, bytes) and BlobViewStruct.is_blob_view_struct(value)): converted_values.append(value) - field_file_ios.append(None) + field_blobs.append(None) continue view_struct = BlobViewStruct.deserialize(value) if blob_view_lookup.resolve_to_null(view_struct): converted_values.append(None) - field_file_ios.append(None) + field_blobs.append(None) else: - descriptor = blob_view_lookup.resolve_descriptor(view_struct) - converted_values.append(descriptor.serialize()) - file_io = blob_view_lookup.resolve_file_io(view_struct) - field_file_ios.append(file_io) + blob = blob_view_lookup.resolve_blob(view_struct) + converted_values.append(blob.to_descriptor().serialize()) + field_blobs.append(blob) column_idx = batch.schema.names.index(field_name) batch = batch.set_column( @@ -159,14 +163,14 @@ def _resolve_view_fields(self, batch, blob_view_lookup): pyarrow.field(field_name, pyarrow.large_binary(), nullable=True), pyarrow.array(converted_values, type=pyarrow.large_binary()), ) - view_file_ios[field_name] = field_file_ios - return batch, view_file_ios + view_blobs[field_name] = field_blobs + return batch, view_blobs # ------------------------------------------------------------------ # Stage 2: BlobData resolution (unified exit) # ------------------------------------------------------------------ - def _resolve_descriptor_fields(self, batch, view_file_ios=None): + def _resolve_descriptor_fields(self, batch, view_blobs=None): if self._blob_as_descriptor: return batch @@ -174,7 +178,10 @@ def _resolve_descriptor_fields(self, batch, view_file_ios=None): if field_name not in batch.schema.names: continue values = [self._normalize_blob_to_bytes(v) for v in batch.column(field_name).to_pylist()] - blobs = [Blob.from_bytes(v, self._table.file_io) for v in values] + blobs = [ + self._descriptor_field_to_blob(value, self._table.file_io) + for value in values + ] if self._blob_parallelism > 1: converted_values = self._table.file_io.read_blobs_concurrent( @@ -189,30 +196,16 @@ def _resolve_descriptor_fields(self, batch, view_file_ios=None): pyarrow.array(converted_values, type=pyarrow.large_binary()), ) - view_file_ios = view_file_ios or {} + view_blobs = view_blobs or {} for field_name in self._view_fields: - field_file_ios = view_file_ios.get(field_name) - if field_name not in batch.schema.names or field_file_ios is None: + blobs = view_blobs.get(field_name) + if field_name not in batch.schema.names or blobs is None: continue - values = [self._normalize_blob_to_bytes(v) for v in batch.column(field_name).to_pylist()] - blobs_by_file_io = {} - converted_values = [] - - for idx, value in enumerate(values): - file_io = field_file_ios[idx] or self._table.file_io - blob = Blob.from_bytes(value, file_io) - if self._blob_parallelism > 1: - converted_values.append(None) - if blob is not None: - blobs_by_file_io.setdefault(file_io, []).append((idx, blob)) - else: - converted_values.append(blob.to_data() if blob else None) - - for file_io, indexed_blobs in blobs_by_file_io.items(): - blobs = [item[1] for item in indexed_blobs] - results = file_io.read_blobs_concurrent(blobs, self._blob_parallelism) - for (idx, _), data in zip(indexed_blobs, results): - converted_values[idx] = data + if self._blob_parallelism > 1: + converted_values = self._table.file_io.read_blobs_concurrent( + blobs, self._blob_parallelism) + else: + converted_values = [blob.to_data() if blob else None for blob in blobs] column_idx = batch.schema.names.index(field_name) batch = batch.set_column( @@ -239,5 +232,19 @@ def _normalize_blob_to_bytes(value): value = bytes(value) return value + @staticmethod + def _descriptor_field_to_blob(value, file_io): + if value is None: + return None + from pypaimon.common.uri_reader import UriReaderFactory + + factory = ( + UriReaderFactory.from_file_io(file_io) if file_io is not None else None) + return Blob.from_descriptor_bytes( + value, + file_io=file_io, + uri_reader_factory=factory, + ) + def close(self): self._inner.close() diff --git a/paimon-python/pypaimon/read/reader/blob_view_read_support.py b/paimon-python/pypaimon/read/reader/blob_view_read_support.py new file mode 100644 index 000000000000..98a7881b872a --- /dev/null +++ b/paimon-python/pypaimon/read/reader/blob_view_read_support.py @@ -0,0 +1,68 @@ +# 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. + +"""Helpers for eager blob-view/descriptor inline conversion on read.""" + +from typing import List + +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.read.reader.iface.record_reader import RecordReader +from pypaimon.schema.data_types import DataField, PyarrowFieldParser + + +def needs_blob_inline_convert(table) -> bool: + view_fields = CoreOptions.blob_view_fields(table.options) + descriptor_fields = CoreOptions.blob_descriptor_fields(table.options) + if descriptor_fields: + # Materialize when blob-as-descriptor=false; otherwise still wrap so + # merge to_iterator()+get_blob() receives descriptor field metadata. + return True + if not view_fields: + return False + if CoreOptions.blob_as_descriptor(table.options): + return True + return CoreOptions.blob_view_resolve_enabled(table.options) + + +def wrap_record_reader_with_blob_inline_convert( + reader: RecordReader, + split_read, + read_fields: List[DataField], +) -> RecordReader: + from pypaimon.read.reader.auth_masking_reader import ( + BatchToRecordReaderAdapter, RecordReaderToBatchAdapter) + from pypaimon.read.reader.blob_descriptor_convert_reader import BlobInlineConvertReader + from pypaimon.read.reader.field_indices import ( + blob_field_indices, descriptor_field_indices_for_table, vector_field_indices) + + schema = PyarrowFieldParser.from_paimon_schema(read_fields) + # Internal round-trip must keep RowKind; default adapter omits _row_kind + # and BatchToRecordReaderAdapter would then emit OffsetRow byte 1 (-U). + batch_reader = RecordReaderToBatchAdapter( + reader, schema, include_row_kind=True) + batch_reader.file_io = split_read.table.file_io + batch_reader.blob_field_indices = blob_field_indices(read_fields) + batch_reader.descriptor_field_indices = descriptor_field_indices_for_table( + split_read.table, read_fields) + batch_reader.vector_field_indices = vector_field_indices(read_fields) + batch_reader = BlobInlineConvertReader( + batch_reader, + split_read.table, + prescan_reader_factory=lambda names: split_read._create_blob_view_prescan_reader(names), + blob_parallelism=split_read._blob_parallelism, + ) + return BatchToRecordReaderAdapter(batch_reader) diff --git a/paimon-python/pypaimon/read/reader/concat_batch_reader.py b/paimon-python/pypaimon/read/reader/concat_batch_reader.py index b1ddd63043e7..059253d60623 100644 --- a/paimon-python/pypaimon/read/reader/concat_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/concat_batch_reader.py @@ -110,11 +110,15 @@ def __init__( class ConcatBatchReader(RecordBatchReader): def __init__(self, reader_suppliers: List[Callable], file_io=None, - blob_field_indices=None, vector_field_indices=None): + blob_field_indices=None, vector_field_indices=None, + descriptor_field_indices=None, + blob_view_lookup=None): self.queue: collections.deque[Callable] = collections.deque(reader_suppliers) self.current_reader: Optional[RecordBatchReader] = None self.file_io = file_io self.blob_field_indices = blob_field_indices + self.descriptor_field_indices = descriptor_field_indices + self.blob_view_lookup = blob_view_lookup self.vector_field_indices = vector_field_indices def read_arrow_batch(self) -> Optional[RecordBatch]: diff --git a/paimon-python/pypaimon/read/reader/deferred_blob_resolve_reader.py b/paimon-python/pypaimon/read/reader/deferred_blob_resolve_reader.py index 9056f00dcb0f..2108d9201243 100644 --- a/paimon-python/pypaimon/read/reader/deferred_blob_resolve_reader.py +++ b/paimon-python/pypaimon/read/reader/deferred_blob_resolve_reader.py @@ -43,6 +43,7 @@ def read_arrow_batch(self) -> Optional[RecordBatch]: batch = self._inner.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) columns = list(batch.columns) fields = list(batch.schema) @@ -52,6 +53,10 @@ def read_arrow_batch(self) -> Optional[RecordBatch]: if column_index < 0: continue values = batch.column(column_index).to_pylist() + # Dedicated .blob files live on the table filesystem. Pass the + # table FileIO so REST tokens apply. Catalog UriReaderFactory.create + # would open an unscoped FileIO; from_file_io is unnecessary here + # because these values are already table-local blob payloads. blobs = [Blob.from_bytes(value, self._file_io) for value in values] if self._blob_parallelism > 1: payloads = self._file_io.read_blobs_concurrent( diff --git a/paimon-python/pypaimon/read/reader/field_indices.py b/paimon-python/pypaimon/read/reader/field_indices.py index 02060a2f519b..78180e707b70 100644 --- a/paimon-python/pypaimon/read/reader/field_indices.py +++ b/paimon-python/pypaimon/read/reader/field_indices.py @@ -29,6 +29,28 @@ def blob_field_indices(fields: List[DataField]) -> Set[int]: } +def descriptor_field_indices( + fields: List[DataField], descriptor_field_names: Iterable[str]) -> Set[int]: + names = set(descriptor_field_names) + if not names: + return set() + return {i for i, f in enumerate(fields) if f.name in names} + + +def descriptor_field_names_for_table(table) -> Set[str]: + from pypaimon.common.options.core_options import CoreOptions + + names = set(CoreOptions.blob_descriptor_fields(table.options)) + if CoreOptions.blob_as_descriptor(table.options): + names |= CoreOptions.blob_view_fields(table.options) + return names + + +def descriptor_field_indices_for_table(table, fields: List[DataField]) -> Set[int]: + return descriptor_field_indices( + fields, descriptor_field_names_for_table(table)) + + def vector_field_indices(fields: List[DataField]) -> Set[int]: return {i for i, f in enumerate(fields) if isinstance(f.type, VectorType)} diff --git a/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py b/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py index fdebbc99cbdf..2100dd93614b 100644 --- a/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py @@ -59,6 +59,7 @@ def read_arrow_batch(self) -> Optional[pa.RecordBatch]: batch = self.reader.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self.reader) if batch.num_rows == 0: return batch filtered = self._filter_batch(batch) @@ -98,6 +99,8 @@ def _filter_batch_by_row(self, batch: pa.RecordBatch) -> Optional[pa.RecordBatch self.file_io, self.blob_field_indices, self.vector_field_indices, + self.descriptor_field_indices, + self.blob_view_lookup, ) selected = [] pos = 0 diff --git a/paimon-python/pypaimon/read/reader/filter_record_reader.py b/paimon-python/pypaimon/read/reader/filter_record_reader.py index 919b9e42e092..11d2407c802f 100644 --- a/paimon-python/pypaimon/read/reader/filter_record_reader.py +++ b/paimon-python/pypaimon/read/reader/filter_record_reader.py @@ -31,11 +31,13 @@ class FilterRecordReader(RecordReader[InternalRow]): def __init__(self, reader: RecordReader[InternalRow], predicate: Predicate): self.reader = reader self.predicate = predicate + self._adopt_blob_metadata(reader) def read_batch(self) -> Optional[RecordIterator[InternalRow]]: iterator = self.reader.read_batch() if iterator is None: return None + self._refresh_blob_view_lookup(self.reader) return FilterRecordIterator(iterator, self.predicate) def close(self) -> None: diff --git a/paimon-python/pypaimon/read/reader/iface/record_batch_reader.py b/paimon-python/pypaimon/read/reader/iface/record_batch_reader.py index 7888f2fab5bf..a87b059a2104 100644 --- a/paimon-python/pypaimon/read/reader/iface/record_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/iface/record_batch_reader.py @@ -36,11 +36,16 @@ class RecordBatchReader(RecordReader): file_io = None blob_field_indices = None + descriptor_field_indices = None + blob_view_lookup = None vector_field_indices = None def _adopt_metadata(self, reader: "RecordBatchReader") -> None: self.file_io = reader.file_io self.blob_field_indices = reader.blob_field_indices + self.descriptor_field_indices = getattr( + reader, 'descriptor_field_indices', None) + self.blob_view_lookup = getattr(reader, 'blob_view_lookup', None) self.vector_field_indices = reader.vector_field_indices @abstractmethod @@ -73,7 +78,8 @@ def read_batch(self) -> Optional[RecordIterator[InternalRow]]: return None return InternalRowWrapperIterator( self._iter_df_rows(df), df.width, self.file_io, - self.blob_field_indices, self.vector_field_indices) + self.blob_field_indices, self.vector_field_indices, + self.descriptor_field_indices, self.blob_view_lookup) @staticmethod def _iter_df_rows(df) -> Iterator[tuple]: @@ -87,12 +93,16 @@ def _iter_df_rows(df) -> Iterator[tuple]: class InternalRowWrapperIterator(RecordIterator[InternalRow]): def __init__(self, iterator: Iterator[tuple], width: int, file_io=None, blob_field_indices=None, - vector_field_indices=None): + vector_field_indices=None, + descriptor_field_indices=None, + blob_view_lookup=None): self._iterator = iterator self._reused_row = OffsetRow(None, 0, width, file_io=file_io, blob_field_indices=blob_field_indices, - vector_field_indices=vector_field_indices) + vector_field_indices=vector_field_indices, + descriptor_field_indices=descriptor_field_indices, + blob_view_lookup=blob_view_lookup) def next(self) -> Optional[InternalRow]: row_tuple = next(self._iterator, None) @@ -114,6 +124,7 @@ def read_arrow_batch(self) -> Optional[RecordBatch]: batch = self._data_reader.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._data_reader) self.batch_pos += batch.num_rows return batch diff --git a/paimon-python/pypaimon/read/reader/iface/record_reader.py b/paimon-python/pypaimon/read/reader/iface/record_reader.py index 2b44955629a7..d777994c6d17 100644 --- a/paimon-python/pypaimon/read/reader/iface/record_reader.py +++ b/paimon-python/pypaimon/read/reader/iface/record_reader.py @@ -39,3 +39,18 @@ def close(self): """ Closes the reader and should release all resources. """ + + def _adopt_blob_metadata(self, reader) -> None: + self.file_io = getattr(reader, 'file_io', None) + self.blob_field_indices = getattr(reader, 'blob_field_indices', None) + self.descriptor_field_indices = getattr(reader, 'descriptor_field_indices', None) + self.blob_view_lookup = getattr(reader, 'blob_view_lookup', None) + self.vector_field_indices = getattr(reader, 'vector_field_indices', None) + + def _refresh_blob_view_lookup(self, reader) -> None: + # BlobInlineConvertReader fills lookup during the first prescan, after + # wrappers have already copied metadata in __init__. Descriptor indices + # are known at construction and must not be overwritten: the inner + # reader may still hold None or unprojected coordinates. + self.blob_view_lookup = getattr( + reader, 'blob_view_lookup', self.blob_view_lookup) diff --git a/paimon-python/pypaimon/read/reader/limited_record_reader.py b/paimon-python/pypaimon/read/reader/limited_record_reader.py index a4eab01986e5..19565c11537d 100644 --- a/paimon-python/pypaimon/read/reader/limited_record_reader.py +++ b/paimon-python/pypaimon/read/reader/limited_record_reader.py @@ -44,6 +44,7 @@ def __init__(self, inner: RecordReader, limit: int): # Public so the iterator can read/write the shared counter without # going through accessor calls per row. self.count = 0 + self._adopt_blob_metadata(inner) def read_batch(self) -> Optional[RecordIterator]: if self.count >= self._limit: @@ -51,6 +52,7 @@ def read_batch(self) -> Optional[RecordIterator]: batch = self._inner.read_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) return _LimitedRecordIterator(batch, self) def close(self) -> None: @@ -96,6 +98,7 @@ def read_arrow_batch(self) -> Optional[RecordBatch]: batch = self._inner.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self._inner) remaining = self._limit - self.count if batch.num_rows > remaining: batch = batch.slice(0, remaining) diff --git a/paimon-python/pypaimon/read/reader/nested_leaf_batch_reader.py b/paimon-python/pypaimon/read/reader/nested_leaf_batch_reader.py index cd849454a719..231ee1513ee5 100644 --- a/paimon-python/pypaimon/read/reader/nested_leaf_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/nested_leaf_batch_reader.py @@ -21,7 +21,8 @@ import pyarrow.compute as pc from pyarrow import RecordBatch -from pypaimon.read.reader.field_indices import blob_field_indices, vector_field_indices +from pypaimon.read.reader.field_indices import ( + blob_field_indices, descriptor_field_indices, vector_field_indices) from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader from pypaimon.schema.data_types import DataField, PyarrowFieldParser @@ -37,7 +38,8 @@ class NestedLeafBatchReader(RecordBatchReader): """ def __init__(self, inner: RecordBatchReader, name_paths: List[List[str]], - output_fields: List[DataField]): + output_fields: List[DataField], + descriptor_field_names=None): if len(name_paths) != len(output_fields): raise ValueError( "name_paths length {} does not match output_fields length {}".format( @@ -47,6 +49,8 @@ def __init__(self, inner: RecordBatchReader, name_paths: List[List[str]], self._schema = PyarrowFieldParser.from_paimon_schema(output_fields) self.file_io = inner.file_io self.blob_field_indices = blob_field_indices(output_fields) + self.descriptor_field_indices = descriptor_field_indices( + output_fields, descriptor_field_names or ()) self.vector_field_indices = vector_field_indices(output_fields) def read_arrow_batch(self) -> Optional[RecordBatch]: diff --git a/paimon-python/pypaimon/read/reader/outer_projection_record_reader.py b/paimon-python/pypaimon/read/reader/outer_projection_record_reader.py index e8bb47509724..90a402b65955 100644 --- a/paimon-python/pypaimon/read/reader/outer_projection_record_reader.py +++ b/paimon-python/pypaimon/read/reader/outer_projection_record_reader.py @@ -45,6 +45,7 @@ def __init__( file_io=None, blob_field_indices=None, vector_field_indices=None, + descriptor_field_indices=None, ): if not name_paths: raise ValueError("name_paths must be non-empty") @@ -65,16 +66,26 @@ def __init__( self._file_io = file_io self._blob_field_indices = project_top_level_field_indices( blob_field_indices, self._specs) + self._descriptor_field_indices = project_top_level_field_indices( + descriptor_field_indices, self._specs) self._vector_field_indices = project_top_level_field_indices( vector_field_indices, self._specs) + self.file_io = self._file_io + self.blob_field_indices = self._blob_field_indices + self.descriptor_field_indices = self._descriptor_field_indices + self.vector_field_indices = self._vector_field_indices + self.blob_view_lookup = getattr(inner, 'blob_view_lookup', None) def read_batch(self) -> Optional[RecordIterator[InternalRow]]: inner_batch = self._inner.read_batch() if inner_batch is None: return None + self._refresh_blob_view_lookup(self._inner) return _OuterProjectionIterator( inner_batch, self._specs, self._flat_arity, self._file_io, - self._blob_field_indices, self._vector_field_indices) + self._blob_field_indices, self._vector_field_indices, + self._descriptor_field_indices, + blob_view_lookup=self.blob_view_lookup) def close(self) -> None: self._inner.close() @@ -91,6 +102,8 @@ def __init__( file_io=None, blob_field_indices=None, vector_field_indices=None, + descriptor_field_indices=None, + blob_view_lookup=None, ): self._inner = inner self._specs = specs @@ -98,7 +111,9 @@ def __init__( self._reused_row = OffsetRow(None, 0, flat_arity, file_io=file_io, blob_field_indices=blob_field_indices, - vector_field_indices=vector_field_indices) + vector_field_indices=vector_field_indices, + descriptor_field_indices=descriptor_field_indices, + blob_view_lookup=blob_view_lookup) def next(self) -> Optional[InternalRow]: inner_row = self._inner.next() diff --git a/paimon-python/pypaimon/read/reader/row_range_filter_record_reader.py b/paimon-python/pypaimon/read/reader/row_range_filter_record_reader.py index d25434f39634..2d7590b4e0c7 100644 --- a/paimon-python/pypaimon/read/reader/row_range_filter_record_reader.py +++ b/paimon-python/pypaimon/read/reader/row_range_filter_record_reader.py @@ -39,6 +39,7 @@ def read_arrow_batch(self) -> Optional[RecordBatch]: batch = self.reader.read_arrow_batch() if batch is None: return None + self._refresh_blob_view_lookup(self.reader) if batch.num_rows == 0: return batch import numpy as np diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index ad31a0103d1d..647d7e5db85c 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -51,10 +51,13 @@ from pypaimon.read.reader.empty_record_reader import EmptyFileRecordReader from pypaimon.read.reader.field_bunch import BlobBunch, DataBunch, FieldBunch, VectorBunch from pypaimon.read.reader.field_indices import ( - blob_field_indices, vector_field_indices) + blob_field_indices, descriptor_field_indices_for_table, + descriptor_field_names_for_table, vector_field_indices) from pypaimon.read.reader.filter_record_reader import FilterRecordReader from pypaimon.read.reader.format_avro_reader import FormatAvroReader from pypaimon.read.reader.blob_descriptor_convert_reader import BlobInlineConvertReader +from pypaimon.read.reader.blob_view_read_support import ( + needs_blob_inline_convert, wrap_record_reader_with_blob_inline_convert) from pypaimon.read.reader.filter_record_batch_reader import FilterRecordBatchReader from pypaimon.read.reader.limited_record_reader import LimitedRecordBatchReader, LimitedRecordReader from pypaimon.read.reader.row_range_filter_record_reader import RowIdFilterRecordBatchReader @@ -185,6 +188,34 @@ def __init__( ) else: self.predicate_for_reader = None + self._blob_view_prescan = False + + def _needs_blob_inline_convert(self) -> bool: + return needs_blob_inline_convert(self.table) + + def _wrap_batch_reader_with_blob_inline_convert( + self, reader: RecordBatchReader) -> RecordBatchReader: + if not self._needs_blob_inline_convert() or self._blob_view_prescan: + return reader + return BlobInlineConvertReader( + reader, + self.table, + prescan_reader_factory=lambda names: self._create_blob_view_prescan_reader(names), + blob_parallelism=self._blob_parallelism, + ) + + def _blob_view_prescan_limit(self) -> Optional[int]: + # Prescan only projects view columns. A predicate/auth filter selects a + # different first-N than LIMIT alone, so do not cap the prescan; the + # outer reader still applies LIMIT after filtering. + if self.predicate is not None: + return None + if getattr(self, '_post_merge_filter', None) is not None: + return None + return self.limit + + def _create_blob_view_prescan_reader(self, field_names: set): + raise NotImplementedError def _compute_nested_path_by_name(self) -> Optional[Dict[str, List[str]]]: if not self.nested_name_paths: @@ -798,8 +829,8 @@ def __init__( row_tracking_enabled: bool, outer_extract_name_paths: Optional[List[List[str]]] = None, outer_flat_read_type: Optional[List[DataField]] = None, - limit: Optional[int] = None): - # Nested-leaf projection is NOT pushed down by name: a leaf path is + limit: Optional[int] = None, + _blob_view_prescan: bool = False): # only valid against the latest schema, while each data file stores # its own (possibly renamed / retyped) sub-fields. Instead the read # widens to the full top-level columns, which the per-file field-id @@ -813,9 +844,25 @@ def __init__( row_tracking_enabled=row_tracking_enabled, nested_name_paths=None, limit=limit) + self._blob_view_prescan = _blob_view_prescan self.outer_extract_name_paths = outer_extract_name_paths self.outer_flat_read_type = outer_flat_read_type + def _create_blob_view_prescan_reader(self, field_names: set): + prescan_fields = [f for f in self.read_fields if f.name in field_names] + if not prescan_fields: + return EmptyRecordBatchReader() + prescan_read = RawFileSplitRead( + table=self.table, + predicate=self.predicate, + read_type=prescan_fields, + split=self.split, + row_tracking_enabled=False, + limit=self._blob_view_prescan_limit(), + _blob_view_prescan=True, + ) + return prescan_read.create_reader() + def raw_reader_supplier(self, file: DataFileMeta, dv_factory: Optional[Callable] = None) -> Optional[RecordReader]: read_fields = self._get_final_read_data_fields() # Check if this is a SlicedSplit to get shard_file_idx_map @@ -866,6 +913,8 @@ def create_reader(self) -> RecordReader: concat_reader = ConcatBatchReader( data_readers, file_io=self.table.file_io, blob_field_indices=blob_field_indices(self.read_fields), + descriptor_field_indices=descriptor_field_indices_for_table( + self.table, self.read_fields), vector_field_indices=vector_field_indices(self.read_fields)) reader = concat_reader if (self.predicate_for_reader @@ -882,7 +931,10 @@ def create_reader(self) -> RecordReader: NestedLeafBatchReader reader = NestedLeafBatchReader( reader, self.outer_extract_name_paths, - self.outer_flat_read_type) + self.outer_flat_read_type, + descriptor_field_names=( + descriptor_field_names_for_table(self.table) + or None)) # A predicate on a projected nested leaf cannot be pushed down: # its leaf path is absent from the widened top-level read fields, # so SplitRead.__init__ dropped it (predicate_for_reader is None). @@ -896,7 +948,7 @@ def create_reader(self) -> RecordReader: reader = FilterRecordBatchReader(reader, trimmed) if self.limit is not None: reader = LimitedRecordBatchReader(reader, self.limit) - return reader + return self._wrap_batch_reader_with_blob_inline_convert(reader) def _all_data_fields_from(self, fields): if self.row_tracking_enabled: @@ -914,7 +966,8 @@ def __init__( row_tracking_enabled: bool, outer_extract_name_paths: Optional[List[List[str]]] = None, outer_flat_read_type: Optional[List[DataField]] = None, - limit: Optional[int] = None): + limit: Optional[int] = None, + _blob_view_prescan: bool = False): self.row_ranges = None if isinstance(split, IndexedSplit): self.row_ranges = split.row_ranges() @@ -933,6 +986,7 @@ def __init__( ) self.outer_extract_name_paths = outer_extract_name_paths self.outer_flat_read_type = outer_flat_read_type + self._blob_view_prescan = _blob_view_prescan # Built once per split-read (value_fields and options are constant # for the object's life), not per section. ``None`` when # ``sequence.field`` is unset, in which case the heap falls back to @@ -1019,6 +1073,65 @@ def _build_merge_function(self): value_field_names=[f.name for f in self.value_fields], ) + def _outer_reapplies_predicate_after_projection(self) -> bool: + return ( + bool(self.outer_extract_name_paths) + and self.predicate is not None + and self.predicate_for_reader is None + and self.outer_flat_read_type is not None + ) + + def _blob_view_prescan_read_type(self, field_names: set): + """View columns plus any ``sequence.field`` needed to merge overlapping files. + + TableRead injects missing sequence fields into the main merge + projection; prescan must do the same or ``builtin_seq_comparator`` + raises ``sequence.field 'ts' not found in value fields ['pic']``. + """ + value_fields = self.read_fields[-self.value_arity:] + prescan_fields = [f for f in value_fields if f.name in field_names] + if not prescan_fields: + return [] + seq_names = self.table.options.sequence_field() + if not seq_names: + return prescan_fields + present = {f.name for f in prescan_fields} + for field in value_fields: + if field.name in seq_names and field.name not in present: + prescan_fields.append(field) + present.add(field.name) + missing = [name for name in seq_names if name not in present] + if missing: + table_fields_by_name = {f.name: f for f in self.table.fields} + for name in missing: + field = table_fields_by_name.get(name) + if field is None: + raise ValueError( + "sequence.field %r not found in table schema" % (name,)) + prescan_fields.append(field) + return prescan_fields + + def _create_blob_view_prescan_reader(self, field_names: set): + prescan_fields = self._blob_view_prescan_read_type(field_names) + if not prescan_fields: + return EmptyRecordBatchReader() + prescan_read = MergeFileSplitRead( + table=self.table, + predicate=self.predicate, + read_type=prescan_fields, + split=self.split, + row_tracking_enabled=False, + limit=self._blob_view_prescan_limit(), + _blob_view_prescan=True, + ) + prescan_read.row_ranges = self.row_ranges + reader = prescan_read.create_reader() + if isinstance(reader, RecordBatchReader): + return reader + from pypaimon.read.reader.auth_masking_reader import RecordReaderToBatchAdapter + schema = PyarrowFieldParser.from_paimon_schema(prescan_fields) + return RecordReaderToBatchAdapter(reader, schema) + def create_reader(self) -> RecordReader: # Create a dict mapping data file name to deletion file reader method self._genarate_deletion_file_readers() @@ -1033,16 +1146,34 @@ def create_reader(self) -> RecordReader: reader = FilterRecordReader(kv_unwrap_reader, self.predicate_for_reader) else: reader = kv_unwrap_reader + value_fields = self.read_fields[-self.value_arity:] + # Apply LIMIT before inline convert so BlobView prescan and the main + # adapter consume the same N rows. Nested-leaf predicates are re-applied + # after outer projection and can drop rows, so keep LIMIT outermost then. + limit_before_convert = ( + self.limit is not None + and not self._outer_reapplies_predicate_after_projection() + ) + if limit_before_convert: + reader = LimitedRecordReader(reader, self.limit) + if self._needs_blob_inline_convert() and not self._blob_view_prescan: + reader = wrap_record_reader_with_blob_inline_convert( + reader, self, value_fields) if self.outer_extract_name_paths: from pypaimon.read.reader.outer_projection_record_reader import \ OuterProjectionRecordReader inner_value_fields = self.read_fields[-self.value_arity:] + inner_descriptor_indices = getattr(reader, 'descriptor_field_indices', None) + if inner_descriptor_indices is None: + inner_descriptor_indices = descriptor_field_indices_for_table( + self.table, inner_value_fields) reader = OuterProjectionRecordReader( reader, [f.name for f in inner_value_fields], self.outer_extract_name_paths, file_io=self.table.file_io, blob_field_indices=blob_field_indices(inner_value_fields), - vector_field_indices=vector_field_indices(inner_value_fields)) + vector_field_indices=vector_field_indices(inner_value_fields), + descriptor_field_indices=inner_descriptor_indices) # A predicate on a projected nested leaf is not pushed down (its leaf # path is absent from the widened-to-full-ROW read fields, so it was # dropped in __init__). Without re-applying it after extraction the @@ -1058,7 +1189,7 @@ def create_reader(self) -> RecordReader: reader, rewrite_predicate_indices( trimmed, self.outer_flat_read_type)) - if self.limit is not None: + if self.limit is not None and not limit_before_convert: reader = LimitedRecordReader(reader, self.limit) return reader @@ -1107,16 +1238,7 @@ def _push_down_predicate(self) -> Optional[Predicate]: def create_reader(self) -> RecordReader: reader = self._create_raw_reader() - - if ((CoreOptions.blob_view_fields(self.table.options) and CoreOptions.blob_view_resolve_enabled( - self.table.options)) - or (not CoreOptions.blob_as_descriptor(self.table.options) - and CoreOptions.blob_descriptor_fields(self.table.options))): - blob_parallelism = self._blob_parallelism - reader = BlobInlineConvertReader( - reader, self.table, - prescan_reader_factory=lambda names: self._create_prescan_reader(names), - blob_parallelism=blob_parallelism) + reader = self._wrap_batch_reader_with_blob_inline_convert(reader) if self._post_filter_after_inline: if self._post_merge_filter is not None: @@ -1172,6 +1294,8 @@ def _create_raw_reader(self) -> RecordReader: merge_reader = ConcatBatchReader( suppliers, file_io=self.table.file_io, blob_field_indices=blob_field_indices(self.read_fields), + descriptor_field_indices=descriptor_field_indices_for_table( + self.table, self.read_fields), vector_field_indices=vector_field_indices(self.read_fields)) if self.predicate_for_reader is not None: reader = FilterRecordBatchReader( @@ -1194,7 +1318,10 @@ def _create_raw_reader(self) -> RecordReader: from pypaimon.read.reader.nested_leaf_batch_reader import \ NestedLeafBatchReader reader = NestedLeafBatchReader( - reader, self.outer_extract_name_paths, self.outer_flat_read_type) + reader, self.outer_extract_name_paths, self.outer_flat_read_type, + descriptor_field_names=( + descriptor_field_names_for_table(self.table) + or None)) if self.limit is not None and not self._post_filter_after_inline: reader = LimitedRecordBatchReader(reader, self.limit) @@ -1246,6 +1373,9 @@ def _selected_local_positions(self, reader_range: Range) -> Optional[List[int]]: for row_id in range(row_range.from_, row_range.to + 1) ] + def _create_blob_view_prescan_reader(self, field_names: set): + return self._create_prescan_reader(field_names) + def _create_prescan_reader(self, field_names): """Create a prescan reader by constructing a new DataEvolutionSplitRead instance that only projects the specified field names. @@ -1259,16 +1389,13 @@ def _create_prescan_reader(self, field_names): if not prescan_fields: return EmptyRecordBatchReader() - # 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 skip_limit else self.limit, + limit=self._blob_view_prescan_limit(), ) prescan_read.row_ranges = self.row_ranges return prescan_read._create_raw_reader() diff --git a/paimon-python/pypaimon/read/table_read.py b/paimon-python/pypaimon/read/table_read.py index 3472666ae5ef..1ca23d1ee3ff 100644 --- a/paimon-python/pypaimon/read/table_read.py +++ b/paimon-python/pypaimon/read/table_read.py @@ -1080,6 +1080,15 @@ def __authed_reader(self, split, auth_result, blob_parallelism=1, if not isinstance(reader, RecordBatchReader): schema = PyarrowFieldParser.from_paimon_schema(effective_read_type) reader = RecordReaderToBatchAdapter(reader, schema, include_row_kind=self.include_row_kind) + if getattr(reader, 'blob_field_indices', None) is None: + from pypaimon.read.reader.field_indices import ( + blob_field_indices, descriptor_field_indices_for_table, + vector_field_indices) + reader.file_io = self.table.file_io + reader.blob_field_indices = blob_field_indices(effective_read_type) + reader.descriptor_field_indices = descriptor_field_indices_for_table( + self.table, effective_read_type) + reader.vector_field_indices = vector_field_indices(effective_read_type) needs_convert_back = True if filter_fn and not embed_filter: diff --git a/paimon-python/pypaimon/table/row/offset_row.py b/paimon-python/pypaimon/table/row/offset_row.py index 4ac8b7dfa13e..4e0255f9692c 100644 --- a/paimon-python/pypaimon/table/row/offset_row.py +++ b/paimon-python/pypaimon/table/row/offset_row.py @@ -25,7 +25,9 @@ class OffsetRow(InternalRow): def __init__(self, row_tuple: Optional[tuple], offset: int, arity: int, file_io=None, blob_field_indices: Optional[Iterable[int]] = None, - vector_field_indices: Optional[Iterable[int]] = None): + vector_field_indices: Optional[Iterable[int]] = None, + descriptor_field_indices: Optional[Iterable[int]] = None, + blob_view_lookup=None): self.row_tuple = row_tuple self.offset = offset self.arity = arity @@ -34,6 +36,11 @@ def __init__(self, row_tuple: Optional[tuple], offset: int, arity: int, self._blob_field_indices: FrozenSet[int] = ( frozenset(blob_field_indices) if blob_field_indices is not None else frozenset() ) + self._descriptor_field_indices: FrozenSet[int] = ( + frozenset(descriptor_field_indices) + if descriptor_field_indices is not None else frozenset() + ) + self._blob_view_lookup = blob_view_lookup self._vector_field_indices: FrozenSet[int] = ( frozenset(vector_field_indices) if vector_field_indices is not None else frozenset() ) @@ -55,12 +62,55 @@ def get_field(self, pos: int): raise IndexError(f"Position {pos} is out of bounds for row arity {self.arity}") return self.row_tuple[self.offset + pos] - def get_blob(self, pos: int): + @staticmethod + def _normalize_blob_bytes(value): + if value is None: + return None + if hasattr(value, 'as_py'): + value = value.as_py() + if isinstance(value, str): + value = value.encode('utf-8') + if isinstance(value, bytearray): + value = bytes(value) + return value + + def _resolve_blob_view_struct(self, view_struct): from pypaimon.table.row.blob import Blob + if self._blob_view_lookup is not None: + if self._blob_view_lookup.resolve_to_null(view_struct): + return None + return self._blob_view_lookup.resolve_blob(view_struct) + return Blob.from_view(view_struct) + + def _blob_from_descriptor_field_bytes(self, raw: bytes): + from pypaimon.table.row.blob import Blob + + return Blob.from_descriptor_bytes( + raw, self._file_io, uri_reader_factory=self._uri_reader_factory()) + + def _uri_reader_factory(self): + if self._file_io is None: + return None + from pypaimon.common.uri_reader import UriReaderFactory + + return UriReaderFactory.from_file_io(self._file_io) + + def get_blob(self, pos: int): + from pypaimon.table.row.blob import Blob, BlobViewStruct + if pos not in self._blob_field_indices: raise TypeError(f"Field at position {pos} is not a BLOB field") - return Blob.from_bytes(self.get_field(pos), self._file_io) + value = self.get_field(pos) + if value is None: + return None + raw = self._normalize_blob_bytes(value) + if raw is not None and BlobViewStruct.is_blob_view_struct(raw): + return self._resolve_blob_view_struct(BlobViewStruct.deserialize(raw)) + if pos in self._descriptor_field_indices: + return self._blob_from_descriptor_field_bytes(raw) + return Blob.from_bytes( + raw, self._file_io, uri_reader_factory=self._uri_reader_factory()) def get_vector(self, pos: int): from pypaimon.table.row.vector import Vector diff --git a/paimon-python/pypaimon/tests/blob_table_test.py b/paimon-python/pypaimon/tests/blob_table_test.py index 1ab02b56ea1f..a6ac33d71755 100755 --- a/paimon-python/pypaimon/tests/blob_table_test.py +++ b/paimon-python/pypaimon/tests/blob_table_test.py @@ -5293,6 +5293,193 @@ def test_legacy_stored_descriptor_fields_keeps_dedicated_blob_layout(self): self.assertEqual(result.num_rows, 1) self.assertEqual(result.column('picture').to_pylist()[0], payload) + def test_blob_view_predicate_and_limit_resolves_filtered_row(self): + """Predicate + LIMIT must not restrict prescan to the unfiltered first-N. + + The matching row can sit past LIMIT in file order; prescan has to + preload that view or convert fails with a missing BlobViewStruct. + """ + from pypaimon import Schema + from pypaimon.table.row.blob import BlobViewStruct + + source_schema = pa.schema([ + ('id', pa.int32()), + ('picture', pa.large_binary()), + ]) + source = Schema.from_pyarrow_schema( + source_schema, + options={ + 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', + } + ) + self.catalog.create_table( + 'test_db.blob_view_pred_limit_source', source, False) + source_table = self.catalog.get_table( + 'test_db.blob_view_pred_limit_source') + + num_rows = 10 + payloads = [f'payload-{i}'.encode() for i in range(num_rows)] + write_builder = source_table.new_batch_write_builder() + writer = write_builder.new_write() + writer.write_arrow(pa.Table.from_pydict({ + 'id': list(range(num_rows)), + 'picture': payloads, + }, schema=source_schema)) + write_builder.new_commit().commit(writer.prepare_commit()) + writer.close() + + picture_field_id = next( + field.id for field in source_table.table_schema.fields + if field.name == 'picture' + ) + view_values = [ + BlobViewStruct( + 'test_db.blob_view_pred_limit_source', picture_field_id, i + ).serialize() + for i in range(num_rows) + ] + + target_schema = pa.schema([ + ('id', pa.int32()), + ('picture', pa.large_binary()), + ]) + target = Schema.from_pyarrow_schema( + target_schema, + options={ + 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', + 'blob-view-field': 'picture', + } + ) + self.catalog.create_table( + 'test_db.blob_view_pred_limit_target', target, False) + target_table = self.catalog.get_table( + 'test_db.blob_view_pred_limit_target') + + target_write_builder = target_table.new_batch_write_builder() + target_writer = target_write_builder.new_write() + target_writer.write_arrow(pa.Table.from_pydict({ + 'id': list(range(num_rows)), + 'picture': view_values, + }, schema=target_schema)) + target_write_builder.new_commit().commit( + target_writer.prepare_commit()) + target_writer.close() + + read_builder = target_table.new_read_builder() + predicate = read_builder.new_predicate_builder().equal("id", 9) + read_builder = read_builder.with_filter(predicate).with_limit(1) + result = read_builder.new_read().to_arrow( + read_builder.new_scan().plan().splits() + ) + self.assertEqual(result.num_rows, 1) + self.assertEqual(result.column('id').to_pylist(), [9]) + self.assertEqual(result.column('picture').to_pylist(), [b'payload-9']) + + def test_blob_view_raw_split_predicate_and_limit_resolves_filtered_row(self): + """Same hole on RawFileSplitRead: view-only prescan plus LIMIT. + + Python schema validation still requires data-evolution for BLOB + tables, so the table is created that way and the read is copied + with data-evolution off to force the append/raw split path. + """ + from unittest import mock + + from pypaimon import Schema + from pypaimon.read.split_read import RawFileSplitRead + from pypaimon.read.table_read import TableRead + from pypaimon.table.row.blob import BlobViewStruct + + source_schema = pa.schema([ + ('id', pa.int32()), + ('picture', pa.large_binary()), + ]) + source = Schema.from_pyarrow_schema( + source_schema, + options={ + 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', + } + ) + self.catalog.create_table( + 'test_db.blob_view_raw_pred_limit_source', source, False) + source_table = self.catalog.get_table( + 'test_db.blob_view_raw_pred_limit_source') + + num_rows = 10 + payloads = [f'raw-payload-{i}'.encode() for i in range(num_rows)] + write_builder = source_table.new_batch_write_builder() + writer = write_builder.new_write() + writer.write_arrow(pa.Table.from_pydict({ + 'id': list(range(num_rows)), + 'picture': payloads, + }, schema=source_schema)) + write_builder.new_commit().commit(writer.prepare_commit()) + writer.close() + + picture_field_id = next( + field.id for field in source_table.table_schema.fields + if field.name == 'picture' + ) + view_values = [ + BlobViewStruct( + 'test_db.blob_view_raw_pred_limit_source', picture_field_id, i + ).serialize() + for i in range(num_rows) + ] + + target_schema = pa.schema([ + ('id', pa.int32()), + ('picture', pa.large_binary()), + ]) + target = Schema.from_pyarrow_schema( + target_schema, + options={ + 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', + 'blob-view-field': 'picture', + } + ) + self.catalog.create_table( + 'test_db.blob_view_raw_pred_limit_target', target, False) + target_table = self.catalog.get_table( + 'test_db.blob_view_raw_pred_limit_target') + + target_write_builder = target_table.new_batch_write_builder() + target_writer = target_write_builder.new_write() + target_writer.write_arrow(pa.Table.from_pydict({ + 'id': list(range(num_rows)), + 'picture': view_values, + }, schema=target_schema)) + target_write_builder.new_commit().commit( + target_writer.prepare_commit()) + target_writer.close() + + raw_table = target_table.copy({'data-evolution.enabled': 'false'}) + self.assertFalse(raw_table.options.data_evolution_enabled()) + + read_builder = raw_table.new_read_builder() + predicate = read_builder.new_predicate_builder().equal("id", 9) + read_builder = read_builder.with_filter(predicate).with_limit(1) + split_types = [] + orig_build = TableRead._build_split_read + + def capturing_build(self, *args, **kwargs): + split_read = orig_build(self, *args, **kwargs) + split_types.append(type(split_read)) + return split_read + + with mock.patch.object(TableRead, '_build_split_read', capturing_build): + result = read_builder.new_read().to_arrow( + read_builder.new_scan().plan().splits() + ) + self.assertIn(RawFileSplitRead, split_types) + self.assertEqual(result.num_rows, 1) + self.assertEqual(result.column('id').to_pylist(), [9]) + self.assertEqual( + result.column('picture').to_pylist(), [b'raw-payload-9']) + class GetBlobTest(unittest.TestCase): diff --git a/paimon-python/pypaimon/tests/blob_test.py b/paimon-python/pypaimon/tests/blob_test.py index 456d8dc2b2a1..57188362292e 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -1652,6 +1652,1198 @@ def test_dedicated_writer_accepts_exact_v1_descriptor_bytes(self): [pa.array([video_bytes + b"x"], type=pa.large_binary())], names=["payload"])) + def test_offset_row_get_blob_uses_table_file_io(self): + from pypaimon.table.row.offset_row import OffsetRow + + data = b"row table blob" + descriptor = BlobDescriptor("file-backed/row.bin", 0, len(data)) + file_io = self._token_aware_file_io(data) + row = OffsetRow( + (descriptor.serialize(),), 0, 1, + file_io=file_io, + blob_field_indices=[0], + descriptor_field_indices=[0], + ) + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobRef) + self.assertEqual(blob.to_data(), data) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + + file_io = self._token_aware_file_io(data) + row = OffsetRow( + (descriptor.serialize(),), 0, 1, + file_io=file_io, + blob_field_indices=[0], + ) + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobRef) + self.assertEqual(blob.to_data(), data) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + + def test_blob_inline_convert_reader_uses_table_file_io(self): + from typing import Optional + + from pyarrow import RecordBatch + + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_descriptor_convert_reader import BlobInlineConvertReader + from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader + + data = b"convert table blob" + descriptor = BlobDescriptor("file-backed/convert.bin", 0, len(data)) + file_io = self._token_aware_file_io(data) + batch = RecordBatch.from_arrays( + [pa.array([descriptor.serialize()], type=pa.large_binary())], + names=["payload"], + ) + + class _InnerReader(RecordBatchReader): + def __init__(self): + self.file_io = file_io + self._batch = batch + self._done = False + + def read_arrow_batch(self) -> Optional[RecordBatch]: + if self._done: + return None + self._done = True + return self._batch + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "payload", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + reader = BlobInlineConvertReader(_InnerReader(), table) + result = reader.read_arrow_batch() + self.assertEqual(result.column("payload").to_pylist(), [data]) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + reader.close() + + def test_offset_row_get_blob_video_frame_descriptor_bytes(self): + from pypaimon.table.row.offset_row import OffsetRow + + data = b"video-frame-payload" + descriptor = VideoFrameDescriptor( + "file-backed/video.mp4", 0, len(data), 2) + file_io = self._token_aware_file_io(data) + row = OffsetRow( + (descriptor.serialize(),), 0, 1, + file_io=file_io, + blob_field_indices=[0], + descriptor_field_indices=[0], + ) + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobRef) + self.assertEqual(blob.to_descriptor(), descriptor) + self.assertEqual(blob.to_data(), data) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + + def test_blob_inline_convert_reader_resolves_video_frame_descriptor(self): + from typing import Optional + + from pyarrow import RecordBatch + + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_descriptor_convert_reader import ( + BlobInlineConvertReader) + from pypaimon.read.reader.iface.record_batch_reader import ( + RecordBatchReader) + + data = b"convert video blob" + descriptor = VideoFrameDescriptor( + "file-backed/convert.mp4", 0, len(data), 2) + file_io = self._token_aware_file_io(data) + batch = RecordBatch.from_arrays( + [pa.array([descriptor.serialize()], type=pa.large_binary())], + names=["payload"], + ) + + class _InnerReader(RecordBatchReader): + def __init__(self): + self.file_io = file_io + self._batch = batch + self._done = False + + def read_arrow_batch(self) -> Optional[RecordBatch]: + if self._done: + return None + self._done = True + return self._batch + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "payload", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + reader = BlobInlineConvertReader(_InnerReader(), table) + result = reader.read_arrow_batch() + self.assertEqual(result.column("payload").to_pylist(), [data]) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + reader.close() + + def test_read_blobs_concurrent_preserves_blob_ref_readers(self): + from unittest.mock import MagicMock + + from pypaimon.common.file_io import FileIO + from pypaimon.common.uri_reader import FileUriReader, UriReader + + shared_uri = "s3://shared/blob" + descriptor = BlobDescriptor(shared_uri, 0, 4) + source_a = MagicMock() + source_b = MagicMock() + source_a.read_ranges_coalesced.return_value = [b"AAAA"] + source_b.read_ranges_coalesced.return_value = [b"BBBB"] + + class MemoryUriReader(UriReader): + def __init__(self): + self.opened = [] + + def new_input_stream(self, uri): + self.opened.append(uri) + return io.BytesIO(b"HTTP") + + http_uri = "https://example.com/blob" + http_reader = MemoryUriReader() + + class _BlobRefSubclass(BlobRef): + def new_input_stream(self): + return io.BytesIO(b"SUBC") + + subclass_reader = FileUriReader(MagicMock()) + blobs = [ + BlobRef(FileUriReader(source_a), descriptor), + BlobRef(FileUriReader(source_b), descriptor), + BlobRef(http_reader, BlobDescriptor(http_uri, 0, 4)), + _BlobRefSubclass(subclass_reader, descriptor), + ] + target_file_io = MagicMock() + + result = FileIO.read_blobs_concurrent(target_file_io, blobs, 4) + + self.assertEqual(result, [b"AAAA", b"BBBB", b"HTTP", b"SUBC"]) + source_a.read_ranges_coalesced.assert_called_once_with( + [(shared_uri, 0, 4)], 4) + source_b.read_ranges_coalesced.assert_called_once_with( + [(shared_uri, 0, 4)], 4) + target_file_io.read_ranges_coalesced.assert_not_called() + subclass_reader.file_io.read_ranges_coalesced.assert_not_called() + self.assertEqual(http_reader.opened, [http_uri]) + + def test_deferred_blob_resolve_reader_uses_table_file_io(self): + from typing import Optional + + from pypaimon.read.reader.deferred_blob_resolve_reader import ( + DeferredBlobResolveReader) + from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader + + data = b"deferred table blob" + descriptor = BlobDescriptor("file-backed/blob.bin", 0, len(data)) + + class FailingFactory: + def create(self, uri): + raise AssertionError( + "dedicated blob resolve should use table FileIO") + + class FileBackedIO: + def __init__(self): + self.uri_reader_factory = FailingFactory() + self.opened_paths = [] + + def new_input_stream(self, path): + self.opened_paths.append(path) + return io.BytesIO(data) + + file_io = FileBackedIO() + batch = pa.RecordBatch.from_arrays( + [pa.array([descriptor.serialize()], type=pa.large_binary())], + names=["payload"], + ) + + class _InnerReader(RecordBatchReader): + def __init__(self): + self._done = False + + def read_arrow_batch(self) -> Optional[pa.RecordBatch]: + if self._done: + return None + self._done = True + return batch + + def close(self): + pass + + reader = DeferredBlobResolveReader(_InnerReader(), file_io, ["payload"]) + result = reader.read_arrow_batch() + self.assertEqual(result.column("payload").to_pylist(), [data]) + self.assertEqual(file_io.opened_paths, [descriptor.uri]) + reader.close() + + def test_offset_row_get_blob_v1_descriptor_bytes(self): + from pypaimon.table.row.offset_row import OffsetRow + + data = b"row-level blob payload" + with tempfile.TemporaryDirectory() as tmp_dir: + blob_path = os.path.join(tmp_dir, "blob.bin") + with open(blob_path, 'wb') as f: + f.write(data) + uri = blob_path.encode('utf-8') + serialized_v1 = ( + bytes([1]) + + struct.pack(' Optional[RecordBatch]: + if self._done: + return None + self._done = True + return self._batch + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "payload", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + inner = _InnerReader() + reader = BlobInlineConvertReader(inner, table) + self.assertEqual(reader.descriptor_field_indices, set()) + + row_iter = reader.read_batch() + self.assertIsNotNone(row_iter) + row = row_iter.next() + self.assertIsNotNone(row) + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobData) + self.assertEqual(blob.to_data(), v1_shaped_inline) + reader.close() + + def test_wrap_record_reader_propagates_blob_metadata_for_get_blob(self): + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_view_read_support import ( + wrap_record_reader_with_blob_inline_convert) + from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.schema.data_types import AtomicType, DataField + from pypaimon.table.row.offset_row import OffsetRow + + data = b"merge-path blob payload" + with tempfile.TemporaryDirectory() as tmp_dir: + blob_path = os.path.join(tmp_dir, "blob.bin") + with open(blob_path, 'wb') as f: + f.write(data) + serialized = BlobDescriptor(blob_path, 0, len(data)).serialize() + file_io = FileIO.get(f"file://{tmp_dir}", {}) + + class _OnceIterator(RecordIterator): + def __init__(self, row): + self._row = row + self._done = False + + def next(self): + if self._done: + return None + self._done = True + return self._row + + class _OnceReader(RecordReader): + def __init__(self, row): + self._row = row + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _OnceIterator(self._row) + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "true", + "blob-descriptor-field": "picture", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + + class _SplitRead: + def __init__(self): + self.table = table + self._blob_parallelism = 1 + + def _create_blob_view_prescan_reader(self, names): + return EmptyRecordBatchReader() + + fields = [DataField(0, "picture", AtomicType("BLOB"))] + inner = _OnceReader(OffsetRow((serialized,), 0, 1)) + wrapped = wrap_record_reader_with_blob_inline_convert( + inner, _SplitRead(), fields) + row_iter = wrapped.read_batch() + self.assertIsNotNone(row_iter) + row = row_iter.next() + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobRef) + self.assertEqual(blob.to_data(), data) + wrapped.close() + + def test_wrap_record_reader_preserves_all_row_kinds(self): + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_view_read_support import ( + wrap_record_reader_with_blob_inline_convert) + from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.schema.data_types import AtomicType, DataField + from pypaimon.table.row.offset_row import OffsetRow + from pypaimon.table.row.row_kind import RowKind + + data = b"row-kind blob payload" + with tempfile.TemporaryDirectory() as tmp_dir: + blob_path = os.path.join(tmp_dir, "blob.bin") + with open(blob_path, 'wb') as f: + f.write(data) + serialized = BlobDescriptor(blob_path, 0, len(data)).serialize() + file_io = FileIO.get(f"file://{tmp_dir}", {}) + + class _KindIterator(RecordIterator): + def __init__(self, rows): + self._rows = list(rows) + + def next(self): + if not self._rows: + return None + return self._rows.pop(0) + + class _KindReader(RecordReader): + def __init__(self, rows): + self._rows = rows + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _KindIterator(self._rows) + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "true", + "blob-descriptor-field": "picture", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + + class _SplitRead: + def __init__(self): + self.table = table + self._blob_parallelism = 1 + + def _create_blob_view_prescan_reader(self, names): + return EmptyRecordBatchReader() + + kinds = ( + RowKind.INSERT, RowKind.UPDATE_BEFORE, + RowKind.UPDATE_AFTER, RowKind.DELETE) + rows = [] + for kind in kinds: + row = OffsetRow((serialized,), 0, 1) + row.set_row_kind_byte(kind.value) + rows.append(row) + wrapped = wrap_record_reader_with_blob_inline_convert( + _KindReader(rows), _SplitRead(), + [DataField(0, "picture", AtomicType("BLOB"))]) + out = [] + batch = wrapped.read_batch() + while batch is not None: + row = batch.next() + while row is not None: + out.append(row.get_row_kind()) + row = batch.next() + batch = wrapped.read_batch() + wrapped.close() + self.assertEqual(list(kinds), out) + + def test_limit_before_wrap_materializes_only_limited_descriptors(self): + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_view_read_support import ( + wrap_record_reader_with_blob_inline_convert) + from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.read.reader.limited_record_reader import LimitedRecordReader + from pypaimon.schema.data_types import AtomicType, DataField + from pypaimon.table.row.offset_row import OffsetRow + + payloads = [b"first-blob-payload", b"second-blob-payload"] + with tempfile.TemporaryDirectory() as tmp_dir: + serialized = [] + for index, payload in enumerate(payloads): + blob_path = os.path.join(tmp_dir, "blob-%d.bin" % index) + with open(blob_path, 'wb') as f: + f.write(payload) + serialized.append( + BlobDescriptor(blob_path, 0, len(payload)).serialize()) + file_io = FileIO.get(f"file://{tmp_dir}", {}) + opened = [] + original_open = file_io.new_input_stream + + def counting_open(path): + opened.append(path) + return original_open(path) + + file_io.new_input_stream = counting_open + + class _Iter(RecordIterator): + def __init__(self, rows): + self._rows = list(rows) + + def next(self): + if not self._rows: + return None + return self._rows.pop(0) + + class _Reader(RecordReader): + def __init__(self, rows): + self._rows = rows + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _Iter(self._rows) + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "picture", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + + class _SplitRead: + def __init__(self): + self.table = table + self._blob_parallelism = 1 + + def _create_blob_view_prescan_reader(self, names): + return EmptyRecordBatchReader() + + rows = [OffsetRow((value,), 0, 1) for value in serialized] + limited = LimitedRecordReader(_Reader(rows), 1) + wrapped = wrap_record_reader_with_blob_inline_convert( + limited, _SplitRead(), + [DataField(0, "picture", AtomicType("BLOB"))]) + batch = wrapped.read_batch() + row = batch.next() + self.assertEqual(row.get_blob(0).to_data(), payloads[0]) + self.assertIsNone(batch.next()) + self.assertIsNone(wrapped.read_batch()) + wrapped.close() + self.assertEqual(1, len(opened)) + + def test_needs_blob_inline_convert_when_blob_as_descriptor(self): + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_view_read_support import ( + needs_blob_inline_convert) + + class _Table: + def __init__(self, options): + self.options = CoreOptions(Options(options)) + + self.assertTrue(needs_blob_inline_convert(_Table({ + "blob-as-descriptor": "true", + "blob-descriptor-field": "picture", + }))) + self.assertTrue(needs_blob_inline_convert(_Table({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "picture", + }))) + self.assertTrue(needs_blob_inline_convert(_Table({ + "blob-as-descriptor": "true", + "blob-view-field": "picture", + }))) + self.assertFalse(needs_blob_inline_convert(_Table({ + "blob-as-descriptor": "true", + }))) + self.assertFalse(needs_blob_inline_convert(_Table({ + "blob.stored-descriptor-fields": "picture", + }))) + + def test_limited_record_reader_keeps_cleared_descriptor_indices(self): + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.auth_masking_reader import ( + BatchToRecordReaderAdapter, RecordReaderToBatchAdapter) + from pypaimon.read.reader.blob_view_read_support import ( + wrap_record_reader_with_blob_inline_convert) + from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.read.reader.limited_record_reader import LimitedRecordReader + from pypaimon.schema.data_types import AtomicType, DataField, PyarrowFieldParser + from pypaimon.table.row.blob import BlobData + from pypaimon.table.row.offset_row import OffsetRow + + data = b"materialized through limit wrapper" + with tempfile.TemporaryDirectory() as tmp_dir: + blob_path = os.path.join(tmp_dir, "blob.bin") + with open(blob_path, 'wb') as f: + f.write(data) + serialized = BlobDescriptor(blob_path, 0, len(data)).serialize() + file_io = FileIO.get(f"file://{tmp_dir}", {}) + + class _OnceIterator(RecordIterator): + def __init__(self, row): + self._row = row + self._done = False + + def next(self): + if self._done: + return None + self._done = True + return self._row + + class _OnceReader(RecordReader): + def __init__(self, row): + self._row = row + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _OnceIterator(self._row) + + def close(self): + pass + + class _CatalogEnvironment: + catalog_loader = None + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "false", + "blob-descriptor-field": "picture", + })) + catalog_environment = _CatalogEnvironment() + + table = _Table() + table.file_io = file_io + + class _SplitRead: + def __init__(self): + self.table = table + self._blob_parallelism = 1 + + def _create_blob_view_prescan_reader(self, names): + return EmptyRecordBatchReader() + + fields = [DataField(0, "picture", AtomicType("BLOB"))] + wrapped = wrap_record_reader_with_blob_inline_convert( + _OnceReader(OffsetRow((serialized,), 0, 1)), _SplitRead(), fields) + limited = LimitedRecordReader(wrapped, 10) + self.assertEqual(limited.descriptor_field_indices, set()) + + schema = PyarrowFieldParser.from_paimon_schema(fields) + batch_reader = RecordReaderToBatchAdapter(limited, schema) + if getattr(batch_reader, 'blob_field_indices', None) is None: + from pypaimon.read.reader.field_indices import ( + blob_field_indices, descriptor_field_indices_for_table, + vector_field_indices) + batch_reader.file_io = file_io + batch_reader.blob_field_indices = blob_field_indices(fields) + batch_reader.descriptor_field_indices = ( + descriptor_field_indices_for_table(table, fields)) + batch_reader.vector_field_indices = vector_field_indices(fields) + reader = BatchToRecordReaderAdapter(batch_reader) + blob = reader.read_batch().next().get_blob(0) + self.assertIsInstance(blob, BlobData) + self.assertEqual(blob.to_data(), data) + reader.close() + + def test_batch_to_record_reader_roundtrip_preserves_get_blob(self): + from pypaimon.read.reader.auth_masking_reader import ( + BatchToRecordReaderAdapter, RecordReaderToBatchAdapter) + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.schema.data_types import AtomicType, DataField, PyarrowFieldParser + from pypaimon.table.row.offset_row import OffsetRow + + data = b"roundtrip blob payload" + with tempfile.TemporaryDirectory() as tmp_dir: + blob_path = os.path.join(tmp_dir, "blob.bin") + with open(blob_path, 'wb') as f: + f.write(data) + serialized = BlobDescriptor(blob_path, 0, len(data)).serialize() + file_io = FileIO.get(f"file://{tmp_dir}", {}) + + class _OnceIterator(RecordIterator): + def __init__(self, row): + self._row = row + self._done = False + + def next(self): + if self._done: + return None + self._done = True + return self._row + + class _OnceReader(RecordReader): + def __init__(self, row): + self._row = row + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _OnceIterator(self._row) + + def close(self): + pass + + fields = [DataField(0, "picture", AtomicType("BLOB"))] + schema = PyarrowFieldParser.from_paimon_schema(fields) + batch_reader = RecordReaderToBatchAdapter( + _OnceReader(OffsetRow((serialized,), 0, 1)), schema) + batch_reader.file_io = file_io + batch_reader.blob_field_indices = {0} + batch_reader.descriptor_field_indices = {0} + first = BatchToRecordReaderAdapter(batch_reader) + second_batch = RecordReaderToBatchAdapter(first, schema) + wrapped = BatchToRecordReaderAdapter(second_batch) + row = wrapped.read_batch().next() + blob = row.get_blob(0) + self.assertIsInstance(blob, BlobRef) + self.assertEqual(blob.to_data(), data) + wrapped.close() + + def test_blob_view_prescan_limit_skips_when_predicate_present(self): + from pypaimon.read.split_read import RawFileSplitRead + + obj = RawFileSplitRead.__new__(RawFileSplitRead) + obj.limit = 1 + obj.predicate = None + obj._post_merge_filter = None + self.assertEqual(1, obj._blob_view_prescan_limit()) + + obj.predicate = object() + self.assertIsNone(obj._blob_view_prescan_limit()) + + obj.predicate = None + obj._post_merge_filter = object() + self.assertIsNone(obj._blob_view_prescan_limit()) + + def test_merge_blob_view_prescan_empty_projection_uses_batch_reader(self): + from pypaimon.read.reader.iface.record_batch_reader import EmptyRecordBatchReader + from pypaimon.read.split_read import MergeFileSplitRead + from pypaimon.schema.data_types import AtomicType, DataField + + obj = MergeFileSplitRead.__new__(MergeFileSplitRead) + obj.read_fields = [ + DataField(0, "_KEY_id", AtomicType("INT")), + DataField(1, "id", AtomicType("INT")), + ] + obj.value_arity = 1 + reader = MergeFileSplitRead._create_blob_view_prescan_reader( + obj, {"picture"}) + self.assertIsInstance(reader, EmptyRecordBatchReader) + self.assertIsNone(reader.read_arrow_batch()) + + def test_merge_blob_view_prescan_keeps_sequence_field(self): + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.split_read import MergeFileSplitRead + from pypaimon.schema.data_types import AtomicType, DataField + + pic = DataField(1, "pic", AtomicType("BLOB")) + ts = DataField(2, "ts", AtomicType("INT")) + obj = MergeFileSplitRead.__new__(MergeFileSplitRead) + obj.read_fields = [ + DataField(0, "_KEY_id", AtomicType("INT")), + DataField(0, "id", AtomicType("INT")), + ts, + pic, + ] + obj.value_arity = 3 + + class _Table: + options = CoreOptions(Options({"sequence.field": "ts"})) + fields = [ + DataField(0, "id", AtomicType("INT")), + ts, + pic, + ] + + obj.table = _Table() + names = [ + f.name for f in + MergeFileSplitRead._blob_view_prescan_read_type(obj, {"pic"}) + ] + self.assertEqual(names, ["pic", "ts"]) + + def test_refresh_blob_view_lookup_does_not_clobber_descriptor_indices(self): + from pypaimon.read.reader.auth_masking_reader import ( + RecordReaderToBatchAdapter) + from pypaimon.read.reader.iface.record_iterator import RecordIterator + from pypaimon.read.reader.iface.record_reader import RecordReader + from pypaimon.read.reader.limited_record_reader import LimitedRecordReader + from pypaimon.read.reader.outer_projection_record_reader import ( + OuterProjectionRecordReader) + from pypaimon.schema.data_types import AtomicType, DataField, PyarrowFieldParser + from pypaimon.table.row.offset_row import OffsetRow + + class _Iter(RecordIterator): + def __init__(self, row): + self._row = row + self._done = False + + def next(self): + if self._done: + return None + self._done = True + return self._row + + class _Reader(RecordReader): + def __init__(self, row): + self._row = row + self._done = False + + def read_batch(self): + if self._done: + return None + self._done = True + return _Iter(self._row) + + def close(self): + pass + + fields = [DataField(0, "picture", AtomicType("BLOB"))] + schema = PyarrowFieldParser.from_paimon_schema(fields) + limited = LimitedRecordReader(_Reader(OffsetRow((b"x",), 0, 1)), 1) + adapter = RecordReaderToBatchAdapter(limited, schema, include_row_kind=True) + adapter.descriptor_field_indices = {0} + self.assertIsNotNone(adapter.read_arrow_batch()) + self.assertEqual(adapter.descriptor_field_indices, {0}) + + inner_row = OffsetRow((0, 1, 2, 3, 4, b"desc"), 0, 6) + projected = OuterProjectionRecordReader( + _Reader(inner_row), + ["a", "b", "c", "d", "e", "payload"], + [["payload"]], + descriptor_field_indices={5}, + ) + self.assertEqual(projected.descriptor_field_indices, {0}) + self.assertIsNotNone(projected.read_batch()) + self.assertEqual(projected.descriptor_field_indices, {0}) + + def test_blob_inline_convert_prescan_empty_projection_reads_main_batch(self): + from typing import Optional + + from pyarrow import RecordBatch + + from pypaimon.common.options import Options + from pypaimon.common.options.core_options import CoreOptions + from pypaimon.read.reader.blob_descriptor_convert_reader import ( + BlobInlineConvertReader) + from pypaimon.read.reader.iface.record_batch_reader import ( + EmptyRecordBatchReader, RecordBatchReader) + + batch = RecordBatch.from_arrays( + [pa.array([1], type=pa.int32())], names=["id"]) + + class _InnerReader(RecordBatchReader): + def __init__(self): + self._done = False + + def read_arrow_batch(self) -> Optional[RecordBatch]: + if self._done: + return None + self._done = True + return batch + + def close(self): + pass + + class _CatalogLoader: + pass + + class _CatalogEnvironment: + catalog_loader = _CatalogLoader() + + class _Table: + options = CoreOptions(Options({ + "blob-as-descriptor": "true", + "blob-view-field": "picture", + })) + catalog_environment = _CatalogEnvironment() + + reader = BlobInlineConvertReader( + _InnerReader(), + _Table(), + prescan_reader_factory=lambda names: EmptyRecordBatchReader(), + ) + result = reader.read_arrow_batch() + self.assertIsNotNone(result) + self.assertEqual(result.column("id").to_pylist(), [1]) + reader.close() + + def test_blob_view_lookup_http_descriptor_uses_http_uri_reader(self): + from unittest.mock import MagicMock + + from pypaimon.common.identifier import Identifier + from pypaimon.common.uri_reader import HttpUriReader, UriReaderFactory + from pypaimon.utils.blob_view_lookup import BlobViewLookup + + class TokenFileIO: + def new_input_stream(self, path): + raise AssertionError( + "HTTP descriptors must not use the table FileIO") + + table_key = "db.source" + view_struct = BlobViewStruct(Identifier.from_string(table_key), 1, 0) + http_uri = "https://example.com/blob.bin" + descriptor = BlobDescriptor(http_uri, 0, 4) + lookup = BlobViewLookup(MagicMock()) + lookup._uri_reader_factory_cache[table_key] = ( + UriReaderFactory.from_file_io(TokenFileIO())) + lookup._store_chunk_results({view_struct: descriptor}, set()) + reader = lookup.resolve_uri_reader(view_struct) + self.assertIsInstance(reader, HttpUriReader) + + def test_blob_view_http_descriptor_materializes_serial_and_parallel(self): + from types import SimpleNamespace + from unittest.mock import MagicMock + + from pypaimon.common.file_io import FileIO + from pypaimon.common.identifier import Identifier + from pypaimon.common.uri_reader import UriReader + from pypaimon.read.reader.blob_descriptor_convert_reader import ( + BlobInlineConvertReader) + + view_struct = BlobViewStruct(Identifier.from_string("db.source"), 1, 0) + descriptor = BlobDescriptor("https://example.com/blob", 0, 4) + batch = pa.RecordBatch.from_arrays( + [pa.array([view_struct.serialize()], type=pa.large_binary())], + names=["picture"], + ) + + class MemoryUriReader(UriReader): + def __init__(self): + self.opened = [] + + def new_input_stream(self, uri): + self.opened.append(uri) + return io.BytesIO(b"DATA") + + class TargetFileIO: + def read_blobs_concurrent(self, blobs, parallelism): + return FileIO.read_blobs_concurrent(self, blobs, parallelism) + + def read_ranges_coalesced(self, ranges, parallelism): + raise AssertionError("HTTP blobs must not use the target FileIO") + + for parallelism in (1, 4): + with self.subTest(parallelism=parallelism): + uri_reader = MemoryUriReader() + lookup = MagicMock() + lookup.resolve_to_null.return_value = False + lookup.resolve_blob.return_value = BlobRef(uri_reader, descriptor) + reader = BlobInlineConvertReader.__new__(BlobInlineConvertReader) + reader._view_fields = {"picture"} + reader._descriptor_fields = set() + reader._blob_as_descriptor = False + reader._blob_parallelism = parallelism + reader._table = SimpleNamespace(file_io=TargetFileIO()) + + descriptor_batch, view_blobs = reader._resolve_view_fields( + batch, lookup) + result = reader._resolve_descriptor_fields( + descriptor_batch, view_blobs) + + self.assertEqual(result.column("picture").to_pylist(), [b"DATA"]) + self.assertEqual(uri_reader.opened, [descriptor.uri]) + lookup.resolve_blob.assert_called_once_with(view_struct) + + def test_internal_row_wrapper_iterator_passes_blob_view_lookup(self): + from unittest.mock import MagicMock + + from pypaimon.read.reader.iface.record_batch_reader import InternalRowWrapperIterator + from pypaimon.table.row.blob import BlobViewStruct + from pypaimon.common.identifier import Identifier + + view_struct = BlobViewStruct(Identifier.from_string("db.source"), 1, 42) + lookup = MagicMock() + lookup.resolve_to_null.return_value = True + iterator = InternalRowWrapperIterator( + iter([(view_struct.serialize(),)]), + 1, + blob_field_indices=[0], + blob_view_lookup=lookup, + ) + row = iterator.next() + self.assertIsNone(row.get_blob(0)) + lookup.resolve_to_null.assert_called_once() + + def _token_aware_file_io(self, data): + class FailingFactory: + def create(self, uri): + raise AssertionError( + "descriptor reads must reuse the table FileIO, not catalog factory") + + class TokenFileIO: + def __init__(self): + self.uri_reader_factory = FailingFactory() + self.opened_paths = [] + + def new_input_stream(self, path): + self.opened_paths.append(path) + return io.BytesIO(data) + + return TokenFileIO() + class BlobEndToEndTest(unittest.TestCase): """End-to-end tests for blob functionality with schema definition, file writing, and reading.""" diff --git a/paimon-python/pypaimon/tests/resolving_file_io_test.py b/paimon-python/pypaimon/tests/resolving_file_io_test.py index c3dec4ebe5c5..ae1f243d8067 100644 --- a/paimon-python/pypaimon/tests/resolving_file_io_test.py +++ b/paimon-python/pypaimon/tests/resolving_file_io_test.py @@ -85,6 +85,40 @@ def test_is_object_store_with_local_warehouse(self): resolving = ResolvingFileIO(opts) self.assertFalse(resolving.is_object_store()) + def test_from_file_io_creates_http_reader(self): + from pypaimon.common.uri_reader import UriReaderFactory + + resolving = ResolvingFileIO(Options({})) + try: + factory = UriReaderFactory.from_file_io(resolving) + reader = factory.create("https://example.com/blob.bin") + self.assertEqual(type(reader).__name__, "HttpUriReader") + finally: + resolving.close() + + def test_from_file_io_reuses_self_for_non_http(self): + import io + + from pypaimon.common.uri_reader import FileUriReader, UriReaderFactory + + resolving = ResolvingFileIO(Options({})) + opened = [] + + def tracking(path): + opened.append(path) + return io.BytesIO(b"ok") + + resolving.new_input_stream = tracking + try: + reader = UriReaderFactory.from_file_io(resolving).create( + "file:///tmp/blob.bin") + self.assertIsInstance(reader, FileUriReader) + self.assertEqual( + reader.new_input_stream("file:///tmp/blob.bin").read(), b"ok") + self.assertEqual(opened, ["file:///tmp/blob.bin"]) + finally: + resolving.close() + class ResolvingFileIOReadWriteTest(unittest.TestCase): """End-to-end read/write tests using ResolvingFileIO with local filesystem.""" diff --git a/paimon-python/pypaimon/tests/rest/rest_token_file_io_test.py b/paimon-python/pypaimon/tests/rest/rest_token_file_io_test.py index dbb919eed28b..731b5307e7a1 100644 --- a/paimon-python/pypaimon/tests/rest/rest_token_file_io_test.py +++ b/paimon-python/pypaimon/tests/rest/rest_token_file_io_test.py @@ -283,6 +283,26 @@ def test_uri_reader_factory_property(self): self.assertTrue(hasattr(uri_reader_factory, 'create'), "uri_reader_factory should support create method") + def test_close_only_closes_instance_uri_reader_factory(self): + with patch.object(RESTTokenFileIO, 'try_to_refresh_token'): + file_io = RESTTokenFileIO( + self.identifier, + self.warehouse_path, + self.catalog_options + ) + + factory = MagicMock() + shared_file_io = MagicMock() + file_io._uri_reader_factory_cache = factory + file_io.file_io = MagicMock(return_value=shared_file_io) + + file_io.close() + file_io.close() + + factory.close.assert_called_once_with() + self.assertIsNone(file_io._uri_reader_factory_cache) + shared_file_io.close.assert_not_called() + def test_filesystem_and_uri_reader_factory_after_serialization(self): with patch.object(RESTTokenFileIO, 'try_to_refresh_token'): original_file_io = RESTTokenFileIO( diff --git a/paimon-python/pypaimon/tests/uri_reader_factory_test.py b/paimon-python/pypaimon/tests/uri_reader_factory_test.py index 5973cdd618c8..4e2948c36236 100644 --- a/paimon-python/pypaimon/tests/uri_reader_factory_test.py +++ b/paimon-python/pypaimon/tests/uri_reader_factory_test.py @@ -18,6 +18,7 @@ import os import tempfile import unittest +import io from pypaimon.common.file_io import FileIO from pypaimon.common.uri_reader import UriReaderFactory, HttpUriReader, FileUriReader, UriReader @@ -120,6 +121,33 @@ def test_cache_size_tracking(self): self.factory.create("http://example.com/another_file.txt") self.assertEqual(self.factory.get_cache_size(), initial_size + 3) + def test_clear_cache_releases_owned_file_ios(self): + self.factory.create(f"file://{self.temp_file}") + self.assertEqual(len(self.factory._owned_file_ios), 1) + self.factory.clear_cache() + self.assertEqual(self.factory.get_cache_size(), 0) + self.assertEqual(self.factory._owned_file_ios, []) + + def test_lru_eviction_keeps_owned_file_ios_until_close(self): + from cachetools import LRUCache + + small_factory = UriReaderFactory({}) + small_factory._readers = LRUCache(1) + small_factory.create(f"file://{self.temp_file}") + self.assertEqual(len(small_factory._owned_file_ios), 1) + small_factory.create("http://example.com/other") + self.assertEqual(len(small_factory._owned_file_ios), 1) + small_factory.clear_cache() + self.assertEqual(len(small_factory._owned_file_ios), 0) + + def test_pickle_resets_reader_cache(self): + import pickle + + self.factory.create(f"file://{self.temp_file}") + restored = pickle.loads(pickle.dumps(self.factory)) + self.assertEqual(restored.get_cache_size(), 0) + self.assertEqual(restored._owned_file_ios, []) + def test_uri_reader_functionality(self): """Test that created URI readers actually work.""" # Test file URI reader @@ -223,6 +251,69 @@ def test_get_file_path_with_file_uri(self): path = UriReader.get_file_path(self.temp_file) self.assertEqual(str(path), self.temp_file) + def test_from_file_io_reuses_provided_file_io_for_non_http(self): + data = b"token-scoped blob" + + class TokenFileIO: + def __init__(self): + self.opened_paths = [] + + def new_input_stream(self, path): + self.opened_paths.append(path) + return io.BytesIO(data) + + file_io = TokenFileIO() + factory = UriReaderFactory.from_file_io(file_io) + self.assertIs(factory, UriReaderFactory.from_file_io(file_io)) + self.assertIsInstance(factory.create("https://example.com/blob.bin"), HttpUriReader) + reader = factory.create("file-backed/blob.bin") + self.assertIsInstance(reader, FileUriReader) + self.assertEqual(reader.new_input_stream("file-backed/blob.bin").read(), data) + self.assertEqual(file_io.opened_paths, ["file-backed/blob.bin"]) + + def test_from_file_io_cache_does_not_pin_file_io(self): + import gc + import weakref + + class TokenFileIO: + def new_input_stream(self, path): + return io.BytesIO(b"x") + + file_io = TokenFileIO() + factory = UriReaderFactory.from_file_io(file_io) + factory.create("file:///tmp/blob.bin") + factory.create("https://example.com/blob.bin") + self.assertIs(factory, UriReaderFactory.from_file_io(file_io)) + self.assertEqual(factory.get_cache_size(), 1) + file_io_ref = weakref.ref(file_io) + del file_io + gc.collect() + self.assertIsNone(file_io_ref()) + + def test_from_file_io_factory_is_pickleable(self): + import pickle + import weakref + + file_io = FileIO.get(self.temp_dir) + factory = UriReaderFactory.from_file_io(file_io) + factory.create(f"file://{self.temp_file}") + factory.create("https://example.com/blob.bin") + # The factory only weakly refs FileIO, so FileIO must be in the pickle + # graph (as it is on a table / ResolvingFileIO). + restored_file_io, restored = pickle.loads(pickle.dumps((file_io, factory))) + self.assertEqual(restored.get_cache_size(), 0) + self.assertIsInstance(restored._provided_file_io, weakref.ref) + self.assertIs(restored._provided_file_io(), restored_file_io) + reader = restored.create(f"file://{self.temp_file}") + self.assertIsInstance(reader, FileUriReader) + stream = reader.new_input_stream(self.temp_file) + try: + self.assertEqual(stream.read().decode('utf-8'), "test content") + finally: + stream.close() + self.assertIsInstance( + restored.create("https://example.com/blob.bin"), HttpUriReader) + if __name__ == '__main__': unittest.main() diff --git a/paimon-python/pypaimon/tests/vector_table_test.py b/paimon-python/pypaimon/tests/vector_table_test.py index 6a4515475ec0..dda0a908d98d 100644 --- a/paimon-python/pypaimon/tests/vector_table_test.py +++ b/paimon-python/pypaimon/tests/vector_table_test.py @@ -24,6 +24,7 @@ from pypaimon import CatalogFactory, Schema from pypaimon.manifest.schema.data_file_meta import DataFileMeta +from pypaimon.table.row.offset_row import OffsetRow from pypaimon.table.row.vector import Vector @@ -70,6 +71,12 @@ def test_empty_vector(self): self.assertEqual(len(v), 0) self.assertEqual(v.to_list(), []) + def test_offset_row_legacy_positional_vector_indices(self): + row = OffsetRow(([1.0, 2.0],), 0, 1, None, None, {0}) + + self.assertEqual(row.get_vector(0).to_list(), [1.0, 2.0]) + self.assertEqual(row._descriptor_field_indices, frozenset()) + class VectorFileDetectionTest(unittest.TestCase): diff --git a/paimon-python/pypaimon/utils/blob_view_lookup.py b/paimon-python/pypaimon/utils/blob_view_lookup.py index 03867b98852d..37d6df5c7eb1 100644 --- a/paimon-python/pypaimon/utils/blob_view_lookup.py +++ b/paimon-python/pypaimon/utils/blob_view_lookup.py @@ -16,10 +16,10 @@ # under the License. from concurrent.futures import ThreadPoolExecutor, as_completed -from typing import Dict, List, Tuple, Set +from typing import Dict, List, Set, Tuple from pypaimon.common.identifier import Identifier -from pypaimon.common.uri_reader import FileUriReader, UriReader +from pypaimon.common.uri_reader import UriReader, UriReaderFactory from pypaimon.common.options.core_options import CoreOptions from pypaimon.table.row.blob import Blob, BlobDescriptor, BlobViewStruct from pypaimon.table.special_fields import SpecialFields @@ -60,6 +60,7 @@ def __init__(self, table): self._table = table self._descriptor_cache: Dict[BlobViewStruct, BlobDescriptor] = {} self._uri_reader_cache: Dict[str, UriReader] = {} + self._uri_reader_factory_cache: Dict[str, UriReaderFactory] = {} self._null_value_cache: Set[BlobViewStruct] = set() def preload(self, view_structs: List[BlobViewStruct]): @@ -79,9 +80,7 @@ def preload(self, view_structs: List[BlobViewStruct]): if len(tasks) <= 1: for plan, range_chunk in tasks: - descriptors, null_values = self._load_descriptor_chunk(plan, range_chunk) - self._descriptor_cache.update(descriptors) - self._null_value_cache.update(null_values) + self._store_chunk_results(*self._load_descriptor_chunk(plan, range_chunk)) return with ThreadPoolExecutor(max_workers=min(_PRELOAD_THREAD_NUM, len(tasks))) as executor: @@ -91,9 +90,7 @@ def preload(self, view_structs: List[BlobViewStruct]): } for future in as_completed(futures): try: - descriptors, null_values = future.result() - self._descriptor_cache.update(descriptors) - self._null_value_cache.update(null_values) + self._store_chunk_results(*future.result()) except Exception as exc: # Cancel remaining futures that have not started yet so a single # failure can abort the rest of the preload work as early as possible. @@ -119,18 +116,14 @@ def resolve_blob(self, view_struct: BlobViewStruct) -> Blob: uri_reader = self.resolve_uri_reader(view_struct) return Blob.from_descriptor(uri_reader, descriptor) - def resolve_file_io(self, view_struct: BlobViewStruct): - uri_reader = self.resolve_uri_reader(view_struct) - if not isinstance(uri_reader, FileUriReader): - raise ValueError( - "Cannot resolve BlobViewStruct {} with parallel blob reads because " - "upstream table {} does not use a file-backed UriReader.".format( - view_struct, view_struct.identifier.get_full_name()) - ) - return uri_reader._file_io - def resolve_uri_reader(self, view_struct: BlobViewStruct) -> UriReader: table_key = view_struct.identifier.get_full_name() + factory = self._uri_reader_factory_cache.get(table_key) + descriptor = self._descriptor_cache.get(view_struct) + if factory is not None and descriptor is not None: + # from_file_io: HTTP(S) stays on HttpUriReader; other URIs reuse + # the upstream table FileIO (REST table token). + return factory.create(descriptor.uri) uri_reader = self._uri_reader_cache.get(table_key) if uri_reader is None: raise ValueError( @@ -139,6 +132,10 @@ def resolve_uri_reader(self, view_struct: BlobViewStruct) -> UriReader: ) return uri_reader + def _store_chunk_results(self, descriptors, null_values): + self._descriptor_cache.update(descriptors) + self._null_value_cache.update(null_values) + def resolve_to_null(self, view_struct: BlobViewStruct) -> bool: if view_struct in self._null_value_cache: return True @@ -162,9 +159,11 @@ def _group_by_table( def _create_table_read_plan(self, table_refs: TableReferences) -> TableReadPlan: upstream_table = self._load_table(table_refs.identifier) - self._uri_reader_cache[table_refs.identifier.get_full_name()] = ( - UriReader.from_file(upstream_table.file_io) - ) + table_key = table_refs.identifier.get_full_name() + self._uri_reader_cache[table_key] = UriReader.from_file( + upstream_table.file_io) + self._uri_reader_factory_cache[table_key] = ( + UriReaderFactory.from_file_io(upstream_table.file_io)) fields: List = [] for field_id in table_refs.references_by_field: