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..20cb76a9f3b6 100644 --- a/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py +++ b/paimon-python/pypaimon/read/reader/blob_descriptor_convert_reader.py @@ -174,7 +174,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 +203,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 +242,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/concat_batch_reader.py b/paimon-python/pypaimon/read/reader/concat_batch_reader.py index 7a0c89f18be8..67f10a476d8f 100644 --- a/paimon-python/pypaimon/read/reader/concat_batch_reader.py +++ b/paimon-python/pypaimon/read/reader/concat_batch_reader.py @@ -53,11 +53,13 @@ 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): 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.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..6f7df1427a85 100644 --- a/paimon-python/pypaimon/read/reader/field_indices.py +++ b/paimon-python/pypaimon/read/reader/field_indices.py @@ -29,6 +29,23 @@ 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_indices_for_table(table, fields: List[DataField]) -> Set[int]: + from pypaimon.common.options.core_options import CoreOptions + + if not CoreOptions.blob_as_descriptor(table.options): + return set() + return descriptor_field_indices( + fields, CoreOptions.blob_descriptor_fields(table.options)) + + 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..00ba0aca5308 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,7 @@ 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, ) 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..b19c358a859b 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,14 @@ class RecordBatchReader(RecordReader): file_io = None blob_field_indices = None + descriptor_field_indices = 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.vector_field_indices = reader.vector_field_indices @abstractmethod @@ -73,7 +76,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) @staticmethod def _iter_df_rows(df) -> Iterator[tuple]: @@ -87,12 +91,14 @@ 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): 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) 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..e7fbae0f82f3 100644 --- a/paimon-python/pypaimon/read/split_read.py +++ b/paimon-python/pypaimon/read/split_read.py @@ -50,7 +50,7 @@ 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, 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 @@ -852,6 +852,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 +868,11 @@ 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=( + CoreOptions.blob_descriptor_fields(self.table.options) + if CoreOptions.blob_as_descriptor(self.table.options) + else 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). @@ -1026,7 +1032,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 @@ -1156,6 +1164,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 +1188,11 @@ 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=( + CoreOptions.blob_descriptor_fields(self.table.options) + if CoreOptions.blob_as_descriptor(self.table.options) + else None)) if self.limit is not None and not self._post_filter_after_inline: reader = LimitedRecordBatchReader(reader, self.limit) diff --git a/paimon-python/pypaimon/table/row/blob.py b/paimon-python/pypaimon/table/row/blob.py index 5a0195c78a16..339895cdd5fd 100644 --- a/paimon-python/pypaimon/table/row/blob.py +++ b/paimon-python/pypaimon/table/row/blob.py @@ -118,6 +118,53 @@ def deserialize(cls, data: bytes) -> '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..9a1d7f76e799 100644 --- a/paimon-python/pypaimon/table/row/offset_row.py +++ b/paimon-python/pypaimon/table/row/offset_row.py @@ -25,6 +25,7 @@ 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, vector_field_indices: Optional[Iterable[int]] = None): self.row_tuple = row_tuple self.offset = offset @@ -34,6 +35,10 @@ 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._vector_field_indices: FrozenSet[int] = ( frozenset(vector_field_indices) if vector_field_indices is not None else frozenset() ) @@ -60,7 +65,10 @@ def get_blob(self, pos: int): 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 pos in self._descriptor_field_indices: + return Blob.from_descriptor_bytes(value, self._file_io) + return Blob.from_bytes(value, 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..bdfb10142f6c 100644 --- a/paimon-python/pypaimon/tests/blob_test.py +++ b/paimon-python/pypaimon/tests/blob_test.py @@ -301,6 +301,115 @@ 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('= 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)