From a7975e59221201675fdd9f6313491f02f2a08a88 Mon Sep 17 00:00:00 2001 From: "wenchao.wu" Date: Mon, 10 Aug 2026 17:06:36 +0800 Subject: [PATCH] [python] BLOB descriptor parsing and read-path compatibility. Introduce explicit descriptor-byte parsing for managed BLOB v1/v2 reads, legacy blob.stored-descriptor-fields fallback, write-path truncation checks, UriReaderFactory lifecycle handling, and row-level descriptor field routing for blob-as-descriptor tables. --- .../pypaimon/common/options/core_options.py | 18 +- paimon-python/pypaimon/common/uri_reader.py | 40 +- .../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 | 1 + .../reader/blob_descriptor_convert_reader.py | 20 +- .../read/reader/blob_view_read_support.py | 53 ++ .../read/reader/concat_batch_reader.py | 6 +- .../pypaimon/read/reader/field_indices.py | 22 + .../read/reader/filter_record_batch_reader.py | 2 + .../read/reader/iface/record_batch_reader.py | 16 +- .../read/reader/nested_leaf_batch_reader.py | 8 +- .../reader/outer_projection_record_reader.py | 10 +- paimon-python/pypaimon/read/split_read.py | 102 +++- paimon-python/pypaimon/table/row/blob.py | 119 +++- .../pypaimon/table/row/offset_row.py | 56 +- paimon-python/pypaimon/tests/blob_test.py | 509 ++++++++++++++++++ .../pypaimon/tests/uri_reader_factory_test.py | 27 + .../pypaimon/write/blob_format_writer.py | 29 +- 20 files changed, 991 insertions(+), 60 deletions(-) create mode 100644 paimon-python/pypaimon/read/reader/blob_view_read_support.py diff --git a/paimon-python/pypaimon/common/options/core_options.py b/paimon-python/pypaimon/common/options/core_options.py index e2cc0ddef2e8..4bb3c1531993 100644 --- a/paimon-python/pypaimon/common/options/core_options.py +++ b/paimon-python/pypaimon/common/options/core_options.py @@ -409,6 +409,15 @@ class CoreOptions: ) ) + BLOB_STORED_DESCRIPTOR_FIELDS: ConfigOption[str] = ( + ConfigOptions.key("blob.stored-descriptor-fields") + .string_type() + .no_default_value() + .with_description( + "Legacy Java option name for blob-descriptor-field." + ) + ) + BLOB_VIEW_FIELD: ConfigOption[str] = ( ConfigOptions.key("blob-view-field") .string_type() @@ -1213,7 +1222,14 @@ def variant_shredding_schema(self) -> Optional[str]: return val def blob_descriptor_fields(self, default=None): - value = self.options.get(CoreOptions.BLOB_DESCRIPTOR_FIELD, default) + value = self.options.get(CoreOptions.BLOB_DESCRIPTOR_FIELD, None) + if value is None: + legacy = self.options.data.get( + CoreOptions.BLOB_STORED_DESCRIPTOR_FIELDS.key()) + if legacy is not None: + value = legacy + else: + value = default return CoreOptions._parse_field_set(value) def blob_view_fields(self, default=None): diff --git a/paimon-python/pypaimon/common/uri_reader.py b/paimon-python/pypaimon/common/uri_reader.py index a417a0c9ae41..4feeb0d83c12 100644 --- a/paimon-python/pypaimon/common/uri_reader.py +++ b/paimon-python/pypaimon/common/uri_reader.py @@ -109,8 +109,13 @@ 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() + self._owned_file_ios = [] + self._closing = False + self._readers = self._new_reader_cache() + + 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 +153,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 +192,13 @@ 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() 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 b315226181a6..3024eecbdf3a 100644 --- a/paimon-python/pypaimon/filesystem/local_file_io.py +++ b/paimon-python/pypaimon/filesystem/local_file_io.py @@ -473,6 +473,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 4423c9559e30..ae8a462d756e 100644 --- a/paimon-python/pypaimon/filesystem/pyarrow_file_io.py +++ b/paimon-python/pypaimon/filesystem/pyarrow_file_io.py @@ -100,6 +100,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..6269f945c2b9 100644 --- a/paimon-python/pypaimon/read/reader/auth_masking_reader.py +++ b/paimon-python/pypaimon/read/reader/auth_masking_reader.py @@ -41,6 +41,7 @@ def __init__(self, inner, schema: pa.Schema, chunk_size: int = 65536, include_ro self._pending_iterator = None self._include_row_kind = include_row_kind self.blob_field_indices = getattr(inner, 'blob_field_indices', None) + self.descriptor_field_indices = getattr(inner, 'descriptor_field_indices', None) self.vector_field_indices = getattr(inner, 'vector_field_indices', None) def read_arrow_batch(self) -> Optional[pa.RecordBatch]: 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..3479210a071d 100644 --- a/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py +++ b/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py @@ -67,12 +67,15 @@ 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() @@ -174,7 +177,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( @@ -200,7 +206,7 @@ def _resolve_descriptor_fields(self, batch, view_file_ios=None): for idx, value in enumerate(values): file_io = field_file_ios[idx] or self._table.file_io - blob = Blob.from_bytes(value, file_io) + blob = self._descriptor_field_to_blob(value, file_io) if self._blob_parallelism > 1: converted_values.append(None) if blob is not None: @@ -239,5 +245,11 @@ 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 + return Blob.from_descriptor_bytes(value, file_io=file_io) + 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..f654a65d6cc4 --- /dev/null +++ b/paimon-python/pypaimon/read/reader/blob_view_read_support.py @@ -0,0 +1,53 @@ +# 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: + return ( + (CoreOptions.blob_view_fields(table.options) + and CoreOptions.blob_view_resolve_enabled(table.options)) + or (not CoreOptions.blob_as_descriptor(table.options) + and CoreOptions.blob_descriptor_fields(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 + + schema = PyarrowFieldParser.from_paimon_schema(read_fields) + batch_reader = RecordReaderToBatchAdapter(reader, schema) + 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 7a0c89f18be8..ca15834f2075 100644 --- a/paimon-python/pypaimon/read/reader/concat_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/concat_batch_reader.py @@ -53,11 +53,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/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..40c41c830123 100644 --- a/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/filter_record_batch_reader.py @@ -98,6 +98,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/iface/record_batch_reader.py b/paimon-python/pypaimon/read/reader/iface/record_batch_reader.py index 7888f2fab5bf..039ce2b64c58 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) 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..0e6eb59f9eec 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,6 +66,8 @@ 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) @@ -74,7 +77,8 @@ def read_batch(self) -> Optional[RecordIterator[InternalRow]]: return None 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) def close(self) -> None: self._inner.close() @@ -91,6 +95,7 @@ def __init__( file_io=None, blob_field_indices=None, vector_field_indices=None, + descriptor_field_indices=None, ): self._inner = inner self._specs = specs @@ -98,7 +103,8 @@ 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) def next(self) -> Optional[InternalRow]: inner_row = self._inner.next() diff --git a/paimon-python/pypaimon/read/split_read.py b/paimon-python/pypaimon/read/split_read.py index 2f6f09665f9b..7d3b148daf7b 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -50,10 +50,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 @@ -182,6 +185,24 @@ 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 _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: @@ -784,8 +805,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 @@ -799,9 +820,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.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 @@ -852,6 +889,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.table.is_primary_key_table and self.predicate_for_reader: @@ -866,7 +905,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). @@ -880,7 +922,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: @@ -898,7 +940,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() @@ -917,6 +960,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 @@ -1003,6 +1047,23 @@ def _build_merge_function(self): value_field_names=[f.name for f in self.value_fields], ) + def _create_blob_view_prescan_reader(self, field_names: set): + 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 EmptyFileRecordReader() + prescan_read = MergeFileSplitRead( + table=self.table, + predicate=self.predicate, + read_type=prescan_fields, + split=self.split, + row_tracking_enabled=False, + limit=self.limit, + _blob_view_prescan=True, + ) + prescan_read.row_ranges = self.row_ranges + return prescan_read.create_reader() + def create_reader(self) -> RecordReader: # Create a dict mapping data file name to deletion file reader method self._genarate_deletion_file_readers() @@ -1017,6 +1078,10 @@ 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:] + 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 @@ -1026,7 +1091,9 @@ def create_reader(self) -> RecordReader: 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=descriptor_field_indices_for_table( + self.table, inner_value_fields)) # 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 @@ -1091,16 +1158,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: @@ -1156,6 +1214,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( @@ -1178,7 +1238,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) @@ -1230,6 +1293,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. diff --git a/paimon-python/pypaimon/table/row/blob.py b/paimon-python/pypaimon/table/row/blob.py index 5a0195c78a16..62d38f55b0e0 100644 --- a/paimon-python/pypaimon/table/row/blob.py +++ b/paimon-python/pypaimon/table/row/blob.py @@ -54,13 +54,13 @@ def version(self) -> int: def serialize(self) -> bytes: uri_bytes = self._uri.encode('utf-8') uri_length = len(uri_bytes) - data = struct.pack(' 1: - data += struct.pack(' 'BlobDescriptor': descriptor._version = version return descriptor + @classmethod + def _try_parse_serialized(cls, data: bytes) -> Optional['BlobDescriptor']: + if not isinstance(data, (bytes, bytearray)): + return None + raw = bytes(data) + if len(raw) < 21: + return None + try: + offset = 0 + version = raw[offset] + offset += 1 + if version < 1 or version > cls.CURRENT_VERSION: + return None + if version > 1: + if offset + 8 > len(raw): + return None + magic = struct.unpack(' len(raw): + return None + uri_length = struct.unpack(' Optional['BlobDescriptor']: + """Best-effort parse when data fully encodes a v1 or v2 descriptor. + + Intended for callers that already know bytes represent a descriptor + (for example :meth:`Blob.from_descriptor_bytes`). The heuristic + :meth:`from_bytes` entry point uses :meth:`is_blob_descriptor` (v2 + magic only) so inline payload bytes are not misclassified as v1 + descriptors. + + Unlike :meth:`is_blob_descriptor` (v2 magic header only), this accepts + v1 descriptors without a magic prefix. Parsing is strict (exact byte + length match) but still heuristic: arbitrary inline blob bytes could + theoretically match. + """ + return cls._try_parse_serialized(data) + @classmethod def is_blob_descriptor(cls, data: bytes) -> bool: if not isinstance(data, (bytes, bytearray)): @@ -386,6 +433,47 @@ def from_file(file_io, file_path: str, offset: int, length: int) -> 'Blob': def from_descriptor(uri_reader: UriReader, descriptor: BlobDescriptor) -> 'Blob': return BlobRef(uri_reader, descriptor) + @staticmethod + def _blob_ref_from_descriptor( + descriptor: 'BlobDescriptor', file_io=None, uri_reader_factory=None, + ) -> 'BlobRef': + if uri_reader_factory is None: + if file_io is None: + raise ValueError("file_io is required to resolve BlobDescriptor bytes") + uri_reader = UriReader.from_file(file_io) + else: + uri_reader = uri_reader_factory.create(descriptor.uri) + return BlobRef(uri_reader, descriptor) + + @staticmethod + def from_descriptor_bytes( + data: Optional[bytes], file_io=None, uri_reader_factory=None, + ) -> Optional['Blob']: + """Build a Blob from bytes known to contain a descriptor. + + Version 1 descriptors have no magic header, so they cannot be + distinguished safely from arbitrary payload bytes. Callers which know + from schema or storage context that a value is a descriptor must use + this method instead of the heuristic :meth:`from_bytes` entry point. + + Malformed or non-descriptor bytes raise :class:`ValueError` (fail-fast). + Trailing padding after the descriptor payload is accepted, matching + Java :meth:`BlobDescriptor.deserialize`. + """ + if data is None: + return None + if not isinstance(data, (bytes, bytearray)): + raise TypeError( + f"Blob.from_descriptor_bytes expects bytes, got {type(data)}") + + try: + descriptor = BlobDescriptor.deserialize(bytes(data)) + except (ValueError, struct.error, UnicodeDecodeError) as exc: + raise ValueError( + "Expected BlobDescriptor bytes, got raw bytes") from exc + return Blob._blob_ref_from_descriptor( + descriptor, file_io=file_io, uri_reader_factory=uri_reader_factory) + @staticmethod def from_view(view_struct: BlobViewStruct) -> 'BlobView': return BlobView(view_struct) @@ -401,20 +489,9 @@ def from_bytes( data = bytes(data) if BlobViewStruct.is_blob_view_struct(data): return Blob.from_view(BlobViewStruct.deserialize(data)) - is_descriptor = BlobDescriptor.is_blob_descriptor(data) - if not allow_blob_data and not is_descriptor: - raise ValueError( - "Expected BlobDescriptor bytes, got raw bytes (allow_blob_data=False)" - ) - if is_descriptor: - descriptor = BlobDescriptor.deserialize(data) - if uri_reader_factory is None: - if file_io is None: - raise ValueError("file_io is required to resolve BlobDescriptor bytes") - uri_reader = UriReader.from_file(file_io) - else: - uri_reader = uri_reader_factory.create(descriptor.uri) - return BlobRef(uri_reader, descriptor) + if BlobDescriptor.is_blob_descriptor(data) or not allow_blob_data: + return Blob.from_descriptor_bytes( + data, file_io=file_io, uri_reader_factory=uri_reader_factory) return BlobData(data) diff --git a/paimon-python/pypaimon/table/row/offset_row.py b/paimon-python/pypaimon/table/row/offset_row.py index 4ac8b7dfa13e..47375b3996f0 100644 --- a/paimon-python/pypaimon/table/row/offset_row.py +++ b/paimon-python/pypaimon/table/row/offset_row.py @@ -25,6 +25,8 @@ class OffsetRow(InternalRow): def __init__(self, row_tuple: Optional[tuple], offset: int, arity: int, file_io=None, blob_field_indices: Optional[Iterable[int]] = None, + descriptor_field_indices: Optional[Iterable[int]] = None, + blob_view_lookup=None, vector_field_indices: Optional[Iterable[int]] = None): self.row_tuple = row_tuple self.offset = offset @@ -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,57 @@ 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 + descriptor = self._blob_view_lookup.resolve_descriptor(view_struct) + uri_reader = self._blob_view_lookup.resolve_uri_reader(view_struct) + return Blob.from_descriptor(uri_reader, descriptor) + return Blob.from_view(view_struct) + + def _blob_from_descriptor_field_bytes(self, raw: bytes): + from pypaimon.table.row.blob import Blob, BlobDescriptor + + if BlobDescriptor.is_blob_descriptor(raw): + return Blob.from_descriptor_bytes(raw, self._file_io) + if BlobDescriptor.parse_if_serialized(raw) is not None: + return Blob.from_descriptor_bytes(raw, self._file_io) + try: + # Accept v1/v2 descriptors with trailing padding (Java deserialize). + return Blob.from_descriptor_bytes(raw, self._file_io) + except ValueError: + # Inline convert may have already materialized payload bytes. + return Blob.from_data(raw) + + 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) def get_vector(self, pos: int): from pypaimon.table.row.vector import Vector diff --git a/paimon-python/pypaimon/tests/blob_test.py b/paimon-python/pypaimon/tests/blob_test.py index f661faebf281..f76b2516d9d1 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -301,6 +301,463 @@ def test_from_bytes_invalid_type_raises(self): with self.assertRaises(TypeError): Blob.from_bytes(12345) + def test_from_descriptor_bytes_rejects_non_descriptor_bytes(self): + with self.assertRaises(ValueError) as ctx: + Blob.from_descriptor_bytes(b"hello blob", file_io=LocalFileIO()) + self.assertIn("Expected BlobDescriptor bytes", str(ctx.exception)) + + def test_from_descriptor_bytes_accepts_v1_descriptor_with_trailing_bytes(self): + import struct + + uri = b"file:///tmp/blob.bin" + 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_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 test_blob_view_struct_roundtrip(self): """Test BlobViewStruct serialization compatibility.""" view_struct = BlobViewStruct("test_db.source_table", 7, 42) @@ -1323,6 +1780,16 @@ def test_blob_descriptor_detection(self): ) random_bytes = b"not-a-descriptor" fake_v1_prefix = b"\x01not-a-descriptor" + # v1-shaped inline payload: passes len/version/uri-length checks but fails + # exact total-length match, so it must not be parsed as a descriptor. + v1_shaped_inline = ( + bytes([1]) + + struct.pack('= 0: + expected_length = descriptor_length stream = blob_value.new_input_stream() try: - chunk = stream.read(self.copy_buffer_size) - while chunk: - crc32 = self._write_with_crc(chunk, crc32) + if expected_length is not None: + crc32 = self._copy_exactly(stream, expected_length, crc32) + else: chunk = stream.read(self.copy_buffer_size) + while chunk: + crc32 = self._write_with_crc(chunk, crc32) + chunk = stream.read(self.copy_buffer_size) finally: stream.close() return blob_pos, self.position - blob_pos, crc32 + def _copy_exactly(self, stream: BinaryIO, length: int, crc32: int) -> int: + remaining = length + while remaining > 0: + to_read = min(self.copy_buffer_size, remaining) + chunk = stream.read(to_read) + if not chunk: + raise EOFError( + "Unexpected EOF while copying BLOB payload: expected %d bytes " + "but source ended %d bytes early." % (length, remaining)) + crc32 = self._write_with_crc(chunk, crc32) + remaining -= len(chunk) + return crc32 + def _write_with_crc(self, data: bytes, crc32: int) -> int: crc32 = crc_backend.crc32(data, crc32) self.output_stream.write(data)