JingsongLi commented on code in PR #9099: URL: https://github.com/apache/paimon/pull/9099#discussion_r3739930047
########## paimon-python/pypaimon/write/hash_index_commit_callback.py: ########## @@ -0,0 +1,33 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +from typing import TYPE_CHECKING + +from pypaimon.write.commit_callback import CommitCallback, CommitCallbackContext + +if TYPE_CHECKING: + from pypaimon.write.row_key_extractor import DynamicBucketRowKeyExtractor + + +class HashIndexSnapshotRefreshCallback(CommitCallback): + """Advance HASH-index snapshot baseline after each successful commit.""" + + def __init__(self, extractor: "DynamicBucketRowKeyExtractor") -> None: + self._extractor = extractor + + def call(self, context: CommitCallbackContext) -> None: + self._extractor.refresh_hash_index_snapshot(context.snapshot) Review Comment: [P1] The HASH baseline also needs to be refreshed when a retry discovers that the commit already succeeded. If the atomic snapshot commit succeeds but the client receives an exception, the retry recognizes the identifier as a duplicate and returns `SuccessResult` before invoking this callback. I reproduced the commit returning successfully with latest snapshot 1 while the writer remained at `base_snapshot_id == 0`; the next checkpoint was then deterministically rejected as stale. Java's commit callback contract has a `retry` path and invokes it for already-committed committables. Please add an equivalent duplicate/retry hook that refreshes from the matched committed snapshot, plus an uncertain-success regression test. ########## paimon-python/pypaimon/blob/primary_key_blob_externalizer.py: ########## @@ -0,0 +1,374 @@ +# 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. + +import uuid +from typing import Callable, List, Optional, Set + +import pyarrow as pa + +from pypaimon.blob.managed_blob_reference_file import ManagedBlobReferenceFile +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.schema.data_types import ( + AtomicType, + DataField, + is_array_blob_type, + is_blob_type, + is_map_blob_type, +) +from pypaimon.table.row.blob import Blob, BlobData, BlobDescriptor +from pypaimon.table.row.generic_row import GenericRow, RowKind +from pypaimon.write.blob_format_writer import BlobFormatWriter + + +class PrimaryKeyBlobExternalizer: + """Externalize primary-key managed BLOB values before they enter the write buffer.""" + + def __init__( + self, + file_io, + value_fields: List[DataField], + managed_blob_fields: Set[str], + new_pack_path: Callable[[], str], + target_file_size: int, + copy_buffer_size: int, + data_file_prefix: str, + descriptor_reader_factory=None): + if target_file_size <= 0: + raise ValueError("Managed BLOB target file size must be positive.") + unknown = set(managed_blob_fields) + self._file_io = file_io + self._field_specs = [] + self._uncommitted_packs: List[str] = [] + self._data_file_prefix = data_file_prefix + for field in value_fields: + if field.name not in managed_blob_fields: + continue + unknown.discard(field.name) + if is_blob_type(field.type): + kind = "scalar" + elif is_array_blob_type(field.type): + kind = "array" + elif is_map_blob_type(field.type): + kind = "map" + else: + raise ValueError( + "Managed BLOB field '%s' must be BLOB, ARRAY<BLOB> or MAP<X, BLOB>, " + "but was %s." % (field.name, field.type)) + self._field_specs.append( + _ManagedBlobField( + field.name, + field.type, + kind, + _ManagedBlobPackWriter( + file_io, + new_pack_path, + target_file_size, + self._uncommitted_packs, + copy_buffer_size, + descriptor_reader_factory, + ), + )) + if unknown: + raise ValueError( + "Managed BLOB fields do not exist in value type: %s." % sorted(unknown)) + + @property + def enabled(self) -> bool: + return bool(self._field_specs) + + def externalize_record_batch(self, batch: pa.RecordBatch) -> pa.RecordBatch: + if not self.enabled: + return batch + columns = list(batch.columns) + names = list(batch.schema.names) + schema_fields = list(batch.schema) + changed = False + try: + for field_spec in self._field_specs: + if field_spec.name not in names: + continue + idx = names.index(field_spec.name) + new_column = field_spec.externalize_column(columns[idx]) + if new_column is not columns[idx]: + columns[idx] = new_column + original_field = schema_fields[idx] + schema_fields[idx] = pa.field( + original_field.name, + new_column.type, + nullable=original_field.nullable, + metadata=original_field.metadata, + ) + changed = True + except Exception: + self.abort() + raise + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=pa.schema(schema_fields)) + + def stage_commit(self) -> List[str]: + """Close packs but retain their ownership until the outer prepare succeeds.""" + try: + for field_spec in self._field_specs: + field_spec.pack_writer.close_current() + return list(self._uncommitted_packs) + except Exception: + self.abort() + raise + + def validate_commit(self, staged_packs: List[str]) -> None: + if self._uncommitted_packs != staged_packs: + raise RuntimeError("Managed BLOB packs changed while commit was staged.") + + def complete_commit(self, staged_packs: List[str]) -> None: + self.validate_commit(staged_packs) + self._uncommitted_packs.clear() + + def prepare_commit(self) -> None: + staged_packs = self.stage_commit() + self.complete_commit(staged_packs) + + def abort(self) -> None: + for field_spec in self._field_specs: + field_spec.pack_writer.abort_current() + for path in self._uncommitted_packs: + self._file_io.delete_quietly(path) + self._uncommitted_packs.clear() + + def close(self) -> None: + self.abort() + factory = self._field_specs[0].pack_writer.descriptor_reader_factory \ + if self._field_specs else None + if getattr(factory, "_blob_descriptor_owned", False): + close = getattr(factory, "close", None) + if callable(close): + close() + + +class _ManagedBlobField: + __slots__ = ("name", "data_type", "kind", "pack_writer") + + def __init__(self, name, data_type, kind, pack_writer): + self.name = name + self.data_type = data_type + self.kind = kind + self.pack_writer = pack_writer + + def externalize_column(self, column: pa.Array) -> pa.Array: + if self.kind == "scalar": + return self._externalize_scalar(column) + if self.kind == "array": + return self._externalize_array(column) + return self._externalize_map(column) + + def _externalize_scalar(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + blob = _to_blob_value( + value, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + out.append(descriptor.serialize()) + changed = True + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_array(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + elements = [] + local_changed = False + for element in value: + if element is None: + elements.append(None) + continue + blob = _to_blob_value( + element, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + elements.append(descriptor.serialize()) + local_changed = True + out.append(elements if local_changed else value) + changed = changed or local_changed + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_map(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + items = [] + local_changed = False + for key, map_value in _map_items(value): + if map_value is None: + items.append((key, None)) + continue + blob = _to_blob_value( + map_value, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + items.append((key, descriptor.serialize())) + local_changed = True + out.append(items if local_changed else value) + changed = changed or local_changed + if not changed: + return column + return pa.array(out, type=column.type) + + +class _ManagedBlobPackWriter: + def __init__( + self, + file_io, + new_pack_path: Callable[[], str], + target_file_size: int, + uncommitted_packs: List[str], + copy_buffer_size: int, + descriptor_reader_factory=None): + self.file_io = file_io + self._new_pack_path = new_pack_path + self._target_file_size = target_file_size + self._uncommitted_packs = uncommitted_packs + self._copy_buffer_size = copy_buffer_size + self.descriptor_reader_factory = descriptor_reader_factory + self._current_path: Optional[str] = None + self._output_stream = None + self._writer: Optional[BlobFormatWriter] = None + self._last_descriptor: Optional[BlobDescriptor] = None + + def write(self, blob: Blob) -> BlobDescriptor: + if self._writer is None: + self._open_current() + self._last_descriptor = None + row = GenericRow( + [blob], + [DataField(0, "blob", AtomicType("BLOB"))], + RowKind.INSERT, + ) + self._writer.add_element(row) Review Comment: [P1] Please reject truncated descriptor sources instead of committing a shortened BLOB. A descriptor-backed `BlobRef` reaches the existing `BlobFormatWriter._write_blob_data` path, which reads until EOF and never verifies the descriptor's declared length. I reproduced a descriptor declaring 10 bytes with only 3 bytes available after its offset: Python produced a new managed descriptor of length 3 and `stage_commit()` succeeded. The Java path preserves the known `BlobRef` length and uses `AbstractBlobElementWriter.copyExactly`, throwing on premature EOF. Please add an exact-copy path for `BlobRef` values with a known length and a regression test for a source truncated after descriptor creation. ########## paimon-python/pypaimon/blob/primary_key_blob_externalizer.py: ########## @@ -0,0 +1,374 @@ +# 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. + +import uuid +from typing import Callable, List, Optional, Set + +import pyarrow as pa + +from pypaimon.blob.managed_blob_reference_file import ManagedBlobReferenceFile +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.schema.data_types import ( + AtomicType, + DataField, + is_array_blob_type, + is_blob_type, + is_map_blob_type, +) +from pypaimon.table.row.blob import Blob, BlobData, BlobDescriptor +from pypaimon.table.row.generic_row import GenericRow, RowKind +from pypaimon.write.blob_format_writer import BlobFormatWriter + + +class PrimaryKeyBlobExternalizer: + """Externalize primary-key managed BLOB values before they enter the write buffer.""" + + def __init__( + self, + file_io, + value_fields: List[DataField], + managed_blob_fields: Set[str], + new_pack_path: Callable[[], str], + target_file_size: int, + copy_buffer_size: int, + data_file_prefix: str, + descriptor_reader_factory=None): + if target_file_size <= 0: + raise ValueError("Managed BLOB target file size must be positive.") + unknown = set(managed_blob_fields) + self._file_io = file_io + self._field_specs = [] + self._uncommitted_packs: List[str] = [] + self._data_file_prefix = data_file_prefix + for field in value_fields: + if field.name not in managed_blob_fields: + continue + unknown.discard(field.name) + if is_blob_type(field.type): + kind = "scalar" + elif is_array_blob_type(field.type): + kind = "array" + elif is_map_blob_type(field.type): + kind = "map" + else: + raise ValueError( + "Managed BLOB field '%s' must be BLOB, ARRAY<BLOB> or MAP<X, BLOB>, " + "but was %s." % (field.name, field.type)) + self._field_specs.append( + _ManagedBlobField( + field.name, + field.type, + kind, + _ManagedBlobPackWriter( + file_io, + new_pack_path, + target_file_size, + self._uncommitted_packs, + copy_buffer_size, + descriptor_reader_factory, + ), + )) + if unknown: + raise ValueError( + "Managed BLOB fields do not exist in value type: %s." % sorted(unknown)) + + @property + def enabled(self) -> bool: + return bool(self._field_specs) + + def externalize_record_batch(self, batch: pa.RecordBatch) -> pa.RecordBatch: + if not self.enabled: + return batch + columns = list(batch.columns) + names = list(batch.schema.names) + schema_fields = list(batch.schema) + changed = False + try: + for field_spec in self._field_specs: + if field_spec.name not in names: + continue + idx = names.index(field_spec.name) + new_column = field_spec.externalize_column(columns[idx]) + if new_column is not columns[idx]: + columns[idx] = new_column + original_field = schema_fields[idx] + schema_fields[idx] = pa.field( + original_field.name, + new_column.type, + nullable=original_field.nullable, + metadata=original_field.metadata, + ) + changed = True + except Exception: + self.abort() + raise + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=pa.schema(schema_fields)) + + def stage_commit(self) -> List[str]: + """Close packs but retain their ownership until the outer prepare succeeds.""" + try: + for field_spec in self._field_specs: + field_spec.pack_writer.close_current() + return list(self._uncommitted_packs) + except Exception: + self.abort() + raise + + def validate_commit(self, staged_packs: List[str]) -> None: + if self._uncommitted_packs != staged_packs: + raise RuntimeError("Managed BLOB packs changed while commit was staged.") + + def complete_commit(self, staged_packs: List[str]) -> None: + self.validate_commit(staged_packs) + self._uncommitted_packs.clear() + + def prepare_commit(self) -> None: + staged_packs = self.stage_commit() + self.complete_commit(staged_packs) + + def abort(self) -> None: + for field_spec in self._field_specs: + field_spec.pack_writer.abort_current() + for path in self._uncommitted_packs: + self._file_io.delete_quietly(path) + self._uncommitted_packs.clear() + + def close(self) -> None: + self.abort() + factory = self._field_specs[0].pack_writer.descriptor_reader_factory \ + if self._field_specs else None + if getattr(factory, "_blob_descriptor_owned", False): + close = getattr(factory, "close", None) + if callable(close): + close() + + +class _ManagedBlobField: + __slots__ = ("name", "data_type", "kind", "pack_writer") + + def __init__(self, name, data_type, kind, pack_writer): + self.name = name + self.data_type = data_type + self.kind = kind + self.pack_writer = pack_writer + + def externalize_column(self, column: pa.Array) -> pa.Array: + if self.kind == "scalar": + return self._externalize_scalar(column) + if self.kind == "array": + return self._externalize_array(column) + return self._externalize_map(column) + + def _externalize_scalar(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + blob = _to_blob_value( + value, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + out.append(descriptor.serialize()) + changed = True + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_array(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + elements = [] + local_changed = False + for element in value: + if element is None: + elements.append(None) + continue + blob = _to_blob_value( + element, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + elements.append(descriptor.serialize()) + local_changed = True + out.append(elements if local_changed else value) + changed = changed or local_changed + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_map(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + items = [] + local_changed = False + for key, map_value in _map_items(value): Review Comment: [P2] Normalize duplicate MAP keys before writing any payloads. PyArrow MAP values may contain duplicate keys, and this loop preserves every entry and externalizes even values that should be overwritten. I reproduced two entries with key `1` producing two output descriptors. The Java externalizer first copies entries into a `LinkedHashMap`, so the last value wins; its tests also cover last-null and binary keys. Please normalize by logical key before externalization (using byte-content equality for binary keys), so only the final value is written to a managed pack. ########## paimon-python/pypaimon/common/options/core_options.py: ########## @@ -1216,10 +1226,16 @@ def blob_descriptor_fields(self, default=None): value = self.options.get(CoreOptions.BLOB_DESCRIPTOR_FIELD, default) return CoreOptions._parse_field_set(value) + def blob_descriptor_source_table(self, default=None): + return self.options.get(CoreOptions.BLOB_DESCRIPTOR_SOURCE_TABLE, default) + def blob_view_fields(self, default=None): value = self.options.get(CoreOptions.BLOB_VIEW_FIELD, default) return CoreOptions._parse_field_set(value) + def blob_inline_fields(self, default=None): + return self.blob_descriptor_fields(default).union(self.blob_view_fields(default)) Review Comment: [P1] Please preserve Java's legacy descriptor-field fallback here. Java declares `blob.stored-descriptor-fields` as a fallback for `blob-descriptor-field`, but `blob_descriptor_fields()` only reads the canonical Python option. Consequently, a Java/legacy schema containing only the fallback key produces an empty `blob_inline_fields()` set and the new PK writer misclassifies an inline descriptor field as managed. That changes the storage and `.blobref` lifecycle semantics across Java/Python compaction. Please implement canonical-key precedence with a fallback to `blob.stored-descriptor-fields`, and add an interoperability test using a Java legacy schema. ########## paimon-python/pypaimon/read/reader/managed_blob_convert_record_reader.py: ########## @@ -0,0 +1,430 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +from typing import Any, Iterable, Optional, Set + +import pyarrow as pa + +from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader +from pypaimon.read.reader.iface.record_iterator import RecordIterator +from pypaimon.read.reader.iface.record_reader import RecordReader +from pypaimon.table.row.blob import Blob, BlobViewStruct +from pypaimon.table.row.internal_row import InternalRow +from pypaimon.table.row.offset_row import OffsetRow + + +class ManagedBlobConvertBatchReader(RecordBatchReader): + """Resolve PK managed and inline BLOB references in Arrow batches.""" + + def __init__( + self, + inner: RecordBatchReader, + file_io, + blob_field_names: Iterable[str], + descriptor_field_names: Optional[Iterable[str]] = None, + view_field_names: Optional[Iterable[str]] = None, + table=None, + blob_as_descriptor: bool = False): + self._inner = inner + self._file_io = file_io + self._blob_field_names = frozenset(blob_field_names) + self._descriptor_field_names = frozenset(descriptor_field_names or ()) + self._view_field_names = frozenset(view_field_names or ()) + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = _BlobViewResolver(table) if self._view_field_names else None + self._adopt_metadata(inner) + + def read_arrow_batch(self) -> Optional[pa.RecordBatch]: + batch = self._inner.read_arrow_batch() + if batch is None: + return None + return convert_primary_key_blob_batch( + batch, + self._file_io, + self._blob_field_names, + self._descriptor_field_names, + self._view_field_names, + self._blob_as_descriptor, + self._view_resolver, + ) + + def close(self) -> None: + try: + self._inner.close() + finally: + if self._view_resolver is not None: + self._view_resolver.close() + + +def convert_managed_blob_batch( + batch: pa.RecordBatch, + file_io, + blob_field_names: Set[str]) -> pa.RecordBatch: + if not blob_field_names: + return batch + columns = [] + changed = False + for field in batch.schema: + column = batch.column(field.name) + if field.name not in blob_field_names: + columns.append(column) + continue + values = [_resolve_blob_payload(value, file_io) for value in column.to_pylist()] + columns.append(pa.array(values, type=field.type)) + changed = True + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=batch.schema) + + +def convert_primary_key_blob_batch( + batch: pa.RecordBatch, + file_io, + managed_blob_field_names: Set[str], + descriptor_field_names: Set[str], + view_field_names: Set[str], + blob_as_descriptor: bool, + view_resolver=None) -> pa.RecordBatch: + """Resolve PK BLOB values after merge, filtering, projection and limit.""" + if view_resolver is not None: + view_resolver.preload_batch(batch, view_field_names) + + columns = [] + changed = False + for field in batch.schema: + column = batch.column(field.name) + values = None + if field.name in managed_blob_field_names and not blob_as_descriptor: + values = [_resolve_blob_payload(value, file_io) for value in column.to_pylist()] + elif field.name in descriptor_field_names and not blob_as_descriptor: + values = [_resolve_descriptor_payload(value, file_io) for value in column.to_pylist()] + elif field.name in view_field_names and view_resolver is not None: + values = [ + view_resolver.resolve(value, blob_as_descriptor) + for value in column.to_pylist() + ] + + if values is None: + columns.append(column) + continue + columns.append(pa.array(values, type=field.type)) + changed = True + + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=batch.schema) + + +class ManagedBlobConvertRecordReader(RecordReader[InternalRow]): + """Resolve PK managed and inline BLOB references in InternalRow batches.""" + + def __init__( + self, + inner: RecordReader[InternalRow], + file_io, + blob_field_indices: Optional[Iterable[int]] = None, + descriptor_field_indices: Optional[Iterable[int]] = None, + view_field_indices: Optional[Iterable[int]] = None, + table=None, + blob_as_descriptor: bool = False): + self._inner = inner + self._file_io = file_io + self._blob_field_indices: Set[int] = ( + frozenset(blob_field_indices) if blob_field_indices is not None else frozenset() + ) + self._descriptor_field_indices: Set[int] = frozenset( + descriptor_field_indices or ()) + self._view_field_indices: Set[int] = frozenset(view_field_indices or ()) + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = _BlobViewResolver(table) if self._view_field_indices else None + + def read_batch(self) -> Optional[RecordIterator[InternalRow]]: + inner_batch = self._inner.read_batch() + if inner_batch is None: + return None + return _ManagedBlobConvertIterator( + inner_batch, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + + def close(self) -> None: + try: + self._inner.close() + finally: + if self._view_resolver is not None: + self._view_resolver.close() + + +class _ManagedBlobConvertIterator(RecordIterator[InternalRow]): + def __init__( + self, + inner: RecordIterator[InternalRow], + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver): + self._inner = inner + self._file_io = file_io + self._blob_field_indices = blob_field_indices + self._descriptor_field_indices = descriptor_field_indices + self._view_field_indices = view_field_indices + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = view_resolver + + def next(self) -> Optional[InternalRow]: + row = self._inner.next() + if row is None: + return None + if not (self._blob_field_indices + or self._descriptor_field_indices + or self._view_field_indices): + return row + if isinstance(row, OffsetRow): + return _convert_offset_row( + row, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + return _ConvertedRow( + row, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + + +def _resolve_blob_payload(value: Any, file_io) -> Any: + if value is None: + return None + if hasattr(value, "as_py"): + value = value.as_py() + if value is None: + return None + if isinstance(value, (bytes, bytearray)): + blob = Blob.from_descriptor_bytes(bytes(value), file_io) + return blob.to_data() if blob is not None else None + if isinstance(value, dict): + return { + key: _resolve_blob_payload(map_value, file_io) + for key, map_value in value.items() + } + if isinstance(value, list): + if value and isinstance(value[0], tuple) and len(value[0]) == 2: + return [ + (key, _resolve_blob_payload(map_value, file_io)) + for key, map_value in value + ] + return [ + _resolve_blob_payload(element, file_io) if element is not None else None + for element in value + ] + if isinstance(value, tuple): + if len(value) == 2 and not isinstance(value[0], (bytes, bytearray)): + key, map_value = value + return key, _resolve_blob_payload(map_value, file_io) + return value + + +def _resolve_descriptor_payload(value: Any, file_io) -> Any: + if value is None: + return None + if hasattr(value, "as_py"): + value = value.as_py() + if value is None: + return None + blob = Blob.from_descriptor_bytes(bytes(value), file_io) + return blob.to_data() if blob is not None else None + + +def _convert_offset_row( + row: OffsetRow, + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver) -> OffsetRow: + values = list(row.row_tuple[row.offset:row.offset + row.arity]) + changed = False + for pos in blob_field_indices if not blob_as_descriptor else (): + if pos >= len(values): + continue + converted = _resolve_blob_payload(values[pos], file_io) + if converted is not values[pos]: + values[pos] = converted + changed = True + for pos in descriptor_field_indices if not blob_as_descriptor else (): + if pos >= len(values): + continue + converted = _resolve_descriptor_payload(values[pos], file_io) + if converted is not values[pos]: + values[pos] = converted + changed = True + if view_resolver is not None: + view_resolver.preload_values( + values[pos] for pos in view_field_indices if pos < len(values)) + for pos in view_field_indices: + if pos >= len(values): + continue + converted = view_resolver.resolve(values[pos], blob_as_descriptor) + if converted is not values[pos]: + values[pos] = converted + changed = True + if not changed: + return row + new_tuple = ( + row.row_tuple[:row.offset] + + tuple(values) + + row.row_tuple[row.offset + row.arity:] + ) + row.replace(new_tuple) + return row + + +class _ConvertedRow(InternalRow): + def __init__( + self, + inner: InternalRow, + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver): + self._inner = inner + self._file_io = file_io + self._blob_field_indices = blob_field_indices + self._descriptor_field_indices = descriptor_field_indices + self._view_field_indices = view_field_indices + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = view_resolver + + def get_field(self, pos: int): + value = self._inner.get_field(pos) + if pos in self._blob_field_indices and not self._blob_as_descriptor: + return _resolve_blob_payload(value, self._file_io) + if pos in self._descriptor_field_indices and not self._blob_as_descriptor: + return _resolve_descriptor_payload(value, self._file_io) + if pos in self._view_field_indices and self._view_resolver is not None: + self._view_resolver.preload_values([value]) + return self._view_resolver.resolve(value, self._blob_as_descriptor) + return value + + def get_blob(self, pos: int): + if not (pos in self._blob_field_indices + or pos in self._descriptor_field_indices + or pos in self._view_field_indices): + return self._inner.get_blob(pos) + value = self.get_field(pos) + if value is None: + return None + if self._blob_as_descriptor: + return Blob.from_descriptor_bytes(value, file_io=self._file_io) + return Blob.from_bytes(value, self._file_io) + + def get_row_kind(self): + return self._inner.get_row_kind() + + def get_vector(self, pos: int): + return self._inner.get_vector(pos) + + def __len__(self) -> int: + return len(self._inner) + + +class _BlobViewResolver: + """Caching facade over BlobViewLookup for the PK output path.""" + + def __init__(self, table): + if table is None: + raise ValueError("table is required to resolve blob-view-field values") + from pypaimon.utils.blob_view_lookup import BlobViewLookup + + self._lookup = BlobViewLookup(table) + self._loaded = set() + + def preload_batch(self, batch: pa.RecordBatch, field_names: Iterable[str]) -> None: + values = ( + value + for field_name in field_names + if field_name in batch.schema.names + for value in batch.column(field_name).to_pylist() + ) + self.preload_values(values) + + def preload_values(self, values: Iterable[Any]) -> None: + structs = [] + pending = set() + for value in values: + raw = _normalize_bytes(value) + if raw is None: + continue + if not BlobViewStruct.is_blob_view_struct(raw): + raise ValueError("Expected BlobViewStruct bytes in blob-view-field value.") + view_struct = BlobViewStruct.deserialize(raw) + if view_struct not in self._loaded and view_struct not in pending: + structs.append(view_struct) + pending.add(view_struct) + if structs: + self._lookup.preload(structs) Review Comment: [P2] Please batch BlobView preloading at split scope rather than invoking the upstream read per row. The OffsetRow path calls `preload_values` for each row; with distinct row IDs, each call builds a catalog/table read plan and scans the upstream table. This turns a normal N-row merge result into O(N) upstream planning/read operations (the Arrow path still does this once per batch). Java pre-scans the complete split, deduplicates all `BlobViewStruct` values, and creates one resolver. Please use the same split-wide model, or at least collect a bounded batch before calling `BlobViewLookup.preload`. ########## paimon-python/pypaimon/blob/primary_key_blob_externalizer.py: ########## @@ -0,0 +1,374 @@ +# 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. + +import uuid +from typing import Callable, List, Optional, Set + +import pyarrow as pa + +from pypaimon.blob.managed_blob_reference_file import ManagedBlobReferenceFile +from pypaimon.common.options.core_options import CoreOptions +from pypaimon.schema.data_types import ( + AtomicType, + DataField, + is_array_blob_type, + is_blob_type, + is_map_blob_type, +) +from pypaimon.table.row.blob import Blob, BlobData, BlobDescriptor +from pypaimon.table.row.generic_row import GenericRow, RowKind +from pypaimon.write.blob_format_writer import BlobFormatWriter + + +class PrimaryKeyBlobExternalizer: + """Externalize primary-key managed BLOB values before they enter the write buffer.""" + + def __init__( + self, + file_io, + value_fields: List[DataField], + managed_blob_fields: Set[str], + new_pack_path: Callable[[], str], + target_file_size: int, + copy_buffer_size: int, + data_file_prefix: str, + descriptor_reader_factory=None): + if target_file_size <= 0: + raise ValueError("Managed BLOB target file size must be positive.") + unknown = set(managed_blob_fields) + self._file_io = file_io + self._field_specs = [] + self._uncommitted_packs: List[str] = [] + self._data_file_prefix = data_file_prefix + for field in value_fields: + if field.name not in managed_blob_fields: + continue + unknown.discard(field.name) + if is_blob_type(field.type): + kind = "scalar" + elif is_array_blob_type(field.type): + kind = "array" + elif is_map_blob_type(field.type): + kind = "map" + else: + raise ValueError( + "Managed BLOB field '%s' must be BLOB, ARRAY<BLOB> or MAP<X, BLOB>, " + "but was %s." % (field.name, field.type)) + self._field_specs.append( + _ManagedBlobField( + field.name, + field.type, + kind, + _ManagedBlobPackWriter( + file_io, + new_pack_path, + target_file_size, + self._uncommitted_packs, + copy_buffer_size, + descriptor_reader_factory, + ), + )) + if unknown: + raise ValueError( + "Managed BLOB fields do not exist in value type: %s." % sorted(unknown)) + + @property + def enabled(self) -> bool: + return bool(self._field_specs) + + def externalize_record_batch(self, batch: pa.RecordBatch) -> pa.RecordBatch: + if not self.enabled: + return batch + columns = list(batch.columns) + names = list(batch.schema.names) + schema_fields = list(batch.schema) + changed = False + try: + for field_spec in self._field_specs: + if field_spec.name not in names: + continue + idx = names.index(field_spec.name) + new_column = field_spec.externalize_column(columns[idx]) + if new_column is not columns[idx]: + columns[idx] = new_column + original_field = schema_fields[idx] + schema_fields[idx] = pa.field( + original_field.name, + new_column.type, + nullable=original_field.nullable, + metadata=original_field.metadata, + ) + changed = True + except Exception: + self.abort() + raise + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=pa.schema(schema_fields)) + + def stage_commit(self) -> List[str]: + """Close packs but retain their ownership until the outer prepare succeeds.""" + try: + for field_spec in self._field_specs: + field_spec.pack_writer.close_current() + return list(self._uncommitted_packs) + except Exception: + self.abort() + raise + + def validate_commit(self, staged_packs: List[str]) -> None: + if self._uncommitted_packs != staged_packs: + raise RuntimeError("Managed BLOB packs changed while commit was staged.") + + def complete_commit(self, staged_packs: List[str]) -> None: + self.validate_commit(staged_packs) + self._uncommitted_packs.clear() + + def prepare_commit(self) -> None: + staged_packs = self.stage_commit() + self.complete_commit(staged_packs) + + def abort(self) -> None: + for field_spec in self._field_specs: + field_spec.pack_writer.abort_current() + for path in self._uncommitted_packs: + self._file_io.delete_quietly(path) + self._uncommitted_packs.clear() + + def close(self) -> None: + self.abort() + factory = self._field_specs[0].pack_writer.descriptor_reader_factory \ + if self._field_specs else None + if getattr(factory, "_blob_descriptor_owned", False): + close = getattr(factory, "close", None) + if callable(close): + close() + + +class _ManagedBlobField: + __slots__ = ("name", "data_type", "kind", "pack_writer") + + def __init__(self, name, data_type, kind, pack_writer): + self.name = name + self.data_type = data_type + self.kind = kind + self.pack_writer = pack_writer + + def externalize_column(self, column: pa.Array) -> pa.Array: + if self.kind == "scalar": + return self._externalize_scalar(column) + if self.kind == "array": + return self._externalize_array(column) + return self._externalize_map(column) + + def _externalize_scalar(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + blob = _to_blob_value( + value, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + out.append(descriptor.serialize()) + changed = True + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_array(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + elements = [] + local_changed = False + for element in value: + if element is None: + elements.append(None) + continue + blob = _to_blob_value( + element, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + elements.append(descriptor.serialize()) + local_changed = True + out.append(elements if local_changed else value) + changed = changed or local_changed + if not changed: + return column + return pa.array(out, type=column.type) + + def _externalize_map(self, column: pa.Array) -> pa.Array: + values = column.to_pylist() + out = [] + changed = False + for value in values: + if value is None: + out.append(None) + continue + items = [] + local_changed = False + for key, map_value in _map_items(value): + if map_value is None: + items.append((key, None)) + continue + blob = _to_blob_value( + map_value, + self.pack_writer.file_io, + self.pack_writer.descriptor_reader_factory, + ) + descriptor = self.pack_writer.write(blob) + items.append((key, descriptor.serialize())) + local_changed = True + out.append(items if local_changed else value) + changed = changed or local_changed + if not changed: + return column + return pa.array(out, type=column.type) + + +class _ManagedBlobPackWriter: + def __init__( + self, + file_io, + new_pack_path: Callable[[], str], + target_file_size: int, + uncommitted_packs: List[str], + copy_buffer_size: int, + descriptor_reader_factory=None): + self.file_io = file_io + self._new_pack_path = new_pack_path + self._target_file_size = target_file_size + self._uncommitted_packs = uncommitted_packs + self._copy_buffer_size = copy_buffer_size + self.descriptor_reader_factory = descriptor_reader_factory + self._current_path: Optional[str] = None + self._output_stream = None + self._writer: Optional[BlobFormatWriter] = None + self._last_descriptor: Optional[BlobDescriptor] = None + + def write(self, blob: Blob) -> BlobDescriptor: + if self._writer is None: + self._open_current() + self._last_descriptor = None + row = GenericRow( + [blob], + [DataField(0, "blob", AtomicType("BLOB"))], + RowKind.INSERT, + ) + self._writer.add_element(row) + descriptor = self._last_descriptor + if descriptor is None: + raise IOError("Managed BLOB writer did not produce a descriptor.") + if self._writer.reach_target_size(self._target_file_size): + self.close_current() + return descriptor + + def _open_current(self) -> None: + self._current_path = self._new_pack_path() + self._uncommitted_packs.append(self._current_path) + self._output_stream = self.file_io.new_output_stream(self._current_path) + + def _consumer(_field_name, descriptor): + self._last_descriptor = descriptor + return False + + self._writer = BlobFormatWriter( + self._output_stream, + blob_consumer=_consumer, + file_path=self._current_path, + copy_buffer_size=self._copy_buffer_size, + ) + + def close_current(self) -> None: + if self._writer is None: + return + try: + self._writer.close() Review Comment: [P2] Keep and close the underlying pack stream when writer close/flush fails. `BlobFormatWriter.close()` flushes before closing the output stream. If flush raises, this `finally` clears `_output_stream`, so the following abort cannot close it. A failing-flush stream reproduces `close_called == False`. Java wraps the pack output in try-with-resources and has a dedicated regression test for this failure. Please retain a local `current_out` and close it in a failure-safe `finally` before clearing these fields, while preserving the original exception. ########## paimon-python/pypaimon/write/writer/key_value_data_writer.py: ########## @@ -63,12 +94,38 @@ def _merge_data(self, existing_data: pa.Table, new_data: pa.Table) -> pa.Table: # N writes incur 1 sort instead of N sorts. return pa.concat_tables([existing_data, new_data]) - def prepare_commit(self) -> List[DataFileMeta]: + def _prepare_commit_data(self) -> Any: if self.pending_data is not None and self.pending_data.num_rows > 0: self._flush_all() - # ``_flush_all`` leaves ``pending_data = None``, so super's - # prepare_commit just returns ``committed_files``. - return super().prepare_commit() + if self._blob_externalizer is not None and self._blob_externalizer.enabled: + return self._blob_externalizer.stage_commit() + return None + + def _complete_prepared_commit_extra(self, extra: Any) -> None: + if self._blob_externalizer is not None and self._blob_externalizer.enabled: + self._blob_externalizer.complete_commit(extra) + + def _validate_prepared_commit_extra(self, extra: Any) -> None: + if self._blob_externalizer is not None and self._blob_externalizer.enabled: + self._blob_externalizer.validate_commit(extra) + + def abort(self): + if self._blob_externalizer is not None and self._blob_externalizer.enabled: + self._blob_externalizer.abort() Review Comment: [P2] Terminal abort currently drops owned descriptor readers without closing them. For `blob-descriptor.source-table` or `blob-descriptor.*`, the externalizer owns a descriptor factory/FileIO, but `abort()` only removes uncommitted packs. `FileStoreWrite.abort()` then clears its bucket writers, while the owned factory is closed only by `PrimaryKeyBlobExternalizer.close()`; there is no remaining release path. Please distinguish staged rollback from terminal writer disposal and close/dispose every bucket writer before removing it, without making retryable rollback close a still-reusable writer. ########## paimon-python/pypaimon/write/table_write.py: ########## @@ -280,6 +312,13 @@ def close(self): self.file_store_write.close() finally: self._release_prepared_indexes() + registered = getattr(self, "_hash_index_commit_callbacks", None) + if registered is not None: + registered.clear() Review Comment: [P2] Clearing this map does not unregister callbacks from live commits. Each `StreamTableCommit` still owns the callback, which retains the extractor and table and continues to invoke the closed writer's extractor on later commits. Repeated writer creation/close with one long-lived commit therefore retains all old writer state. Before clearing registrations/detaching the writer, remove each callback from its corresponding commit (or have the builder perform the reverse cleanup). ########## paimon-python/pypaimon/read/reader/managed_blob_convert_record_reader.py: ########## @@ -0,0 +1,430 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +from typing import Any, Iterable, Optional, Set + +import pyarrow as pa + +from pypaimon.read.reader.iface.record_batch_reader import RecordBatchReader +from pypaimon.read.reader.iface.record_iterator import RecordIterator +from pypaimon.read.reader.iface.record_reader import RecordReader +from pypaimon.table.row.blob import Blob, BlobViewStruct +from pypaimon.table.row.internal_row import InternalRow +from pypaimon.table.row.offset_row import OffsetRow + + +class ManagedBlobConvertBatchReader(RecordBatchReader): + """Resolve PK managed and inline BLOB references in Arrow batches.""" + + def __init__( + self, + inner: RecordBatchReader, + file_io, + blob_field_names: Iterable[str], + descriptor_field_names: Optional[Iterable[str]] = None, + view_field_names: Optional[Iterable[str]] = None, + table=None, + blob_as_descriptor: bool = False): + self._inner = inner + self._file_io = file_io + self._blob_field_names = frozenset(blob_field_names) + self._descriptor_field_names = frozenset(descriptor_field_names or ()) + self._view_field_names = frozenset(view_field_names or ()) + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = _BlobViewResolver(table) if self._view_field_names else None + self._adopt_metadata(inner) + + def read_arrow_batch(self) -> Optional[pa.RecordBatch]: + batch = self._inner.read_arrow_batch() + if batch is None: + return None + return convert_primary_key_blob_batch( + batch, + self._file_io, + self._blob_field_names, + self._descriptor_field_names, + self._view_field_names, + self._blob_as_descriptor, + self._view_resolver, + ) + + def close(self) -> None: + try: + self._inner.close() + finally: + if self._view_resolver is not None: + self._view_resolver.close() + + +def convert_managed_blob_batch( + batch: pa.RecordBatch, + file_io, + blob_field_names: Set[str]) -> pa.RecordBatch: + if not blob_field_names: + return batch + columns = [] + changed = False + for field in batch.schema: + column = batch.column(field.name) + if field.name not in blob_field_names: + columns.append(column) + continue + values = [_resolve_blob_payload(value, file_io) for value in column.to_pylist()] + columns.append(pa.array(values, type=field.type)) + changed = True + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=batch.schema) + + +def convert_primary_key_blob_batch( + batch: pa.RecordBatch, + file_io, + managed_blob_field_names: Set[str], + descriptor_field_names: Set[str], + view_field_names: Set[str], + blob_as_descriptor: bool, + view_resolver=None) -> pa.RecordBatch: + """Resolve PK BLOB values after merge, filtering, projection and limit.""" + if view_resolver is not None: + view_resolver.preload_batch(batch, view_field_names) + + columns = [] + changed = False + for field in batch.schema: + column = batch.column(field.name) + values = None + if field.name in managed_blob_field_names and not blob_as_descriptor: + values = [_resolve_blob_payload(value, file_io) for value in column.to_pylist()] + elif field.name in descriptor_field_names and not blob_as_descriptor: + values = [_resolve_descriptor_payload(value, file_io) for value in column.to_pylist()] + elif field.name in view_field_names and view_resolver is not None: + values = [ + view_resolver.resolve(value, blob_as_descriptor) + for value in column.to_pylist() + ] + + if values is None: + columns.append(column) + continue + columns.append(pa.array(values, type=field.type)) + changed = True + + if not changed: + return batch + return pa.RecordBatch.from_arrays(columns, schema=batch.schema) + + +class ManagedBlobConvertRecordReader(RecordReader[InternalRow]): + """Resolve PK managed and inline BLOB references in InternalRow batches.""" + + def __init__( + self, + inner: RecordReader[InternalRow], + file_io, + blob_field_indices: Optional[Iterable[int]] = None, + descriptor_field_indices: Optional[Iterable[int]] = None, + view_field_indices: Optional[Iterable[int]] = None, + table=None, + blob_as_descriptor: bool = False): + self._inner = inner + self._file_io = file_io + self._blob_field_indices: Set[int] = ( + frozenset(blob_field_indices) if blob_field_indices is not None else frozenset() + ) + self._descriptor_field_indices: Set[int] = frozenset( + descriptor_field_indices or ()) + self._view_field_indices: Set[int] = frozenset(view_field_indices or ()) + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = _BlobViewResolver(table) if self._view_field_indices else None + + def read_batch(self) -> Optional[RecordIterator[InternalRow]]: + inner_batch = self._inner.read_batch() + if inner_batch is None: + return None + return _ManagedBlobConvertIterator( + inner_batch, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + + def close(self) -> None: + try: + self._inner.close() + finally: + if self._view_resolver is not None: + self._view_resolver.close() + + +class _ManagedBlobConvertIterator(RecordIterator[InternalRow]): + def __init__( + self, + inner: RecordIterator[InternalRow], + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver): + self._inner = inner + self._file_io = file_io + self._blob_field_indices = blob_field_indices + self._descriptor_field_indices = descriptor_field_indices + self._view_field_indices = view_field_indices + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = view_resolver + + def next(self) -> Optional[InternalRow]: + row = self._inner.next() + if row is None: + return None + if not (self._blob_field_indices + or self._descriptor_field_indices + or self._view_field_indices): + return row + if isinstance(row, OffsetRow): + return _convert_offset_row( + row, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + return _ConvertedRow( + row, + self._file_io, + self._blob_field_indices, + self._descriptor_field_indices, + self._view_field_indices, + self._blob_as_descriptor, + self._view_resolver, + ) + + +def _resolve_blob_payload(value: Any, file_io) -> Any: + if value is None: + return None + if hasattr(value, "as_py"): + value = value.as_py() + if value is None: + return None + if isinstance(value, (bytes, bytearray)): + blob = Blob.from_descriptor_bytes(bytes(value), file_io) + return blob.to_data() if blob is not None else None + if isinstance(value, dict): + return { + key: _resolve_blob_payload(map_value, file_io) + for key, map_value in value.items() + } + if isinstance(value, list): + if value and isinstance(value[0], tuple) and len(value[0]) == 2: + return [ + (key, _resolve_blob_payload(map_value, file_io)) + for key, map_value in value + ] + return [ + _resolve_blob_payload(element, file_io) if element is not None else None + for element in value + ] + if isinstance(value, tuple): + if len(value) == 2 and not isinstance(value[0], (bytes, bytearray)): + key, map_value = value + return key, _resolve_blob_payload(map_value, file_io) + return value + + +def _resolve_descriptor_payload(value: Any, file_io) -> Any: + if value is None: + return None + if hasattr(value, "as_py"): + value = value.as_py() + if value is None: + return None + blob = Blob.from_descriptor_bytes(bytes(value), file_io) + return blob.to_data() if blob is not None else None + + +def _convert_offset_row( + row: OffsetRow, + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver) -> OffsetRow: + values = list(row.row_tuple[row.offset:row.offset + row.arity]) + changed = False + for pos in blob_field_indices if not blob_as_descriptor else (): + if pos >= len(values): + continue + converted = _resolve_blob_payload(values[pos], file_io) + if converted is not values[pos]: + values[pos] = converted + changed = True + for pos in descriptor_field_indices if not blob_as_descriptor else (): + if pos >= len(values): + continue + converted = _resolve_descriptor_payload(values[pos], file_io) + if converted is not values[pos]: + values[pos] = converted + changed = True + if view_resolver is not None: + view_resolver.preload_values( + values[pos] for pos in view_field_indices if pos < len(values)) + for pos in view_field_indices: + if pos >= len(values): + continue + converted = view_resolver.resolve(values[pos], blob_as_descriptor) + if converted is not values[pos]: + values[pos] = converted + changed = True + if not changed: + return row + new_tuple = ( + row.row_tuple[:row.offset] + + tuple(values) + + row.row_tuple[row.offset + row.arity:] + ) + row.replace(new_tuple) + return row + + +class _ConvertedRow(InternalRow): + def __init__( + self, + inner: InternalRow, + file_io, + blob_field_indices: Set[int], + descriptor_field_indices: Set[int], + view_field_indices: Set[int], + blob_as_descriptor: bool, + view_resolver): + self._inner = inner + self._file_io = file_io + self._blob_field_indices = blob_field_indices + self._descriptor_field_indices = descriptor_field_indices + self._view_field_indices = view_field_indices + self._blob_as_descriptor = blob_as_descriptor + self._view_resolver = view_resolver + + def get_field(self, pos: int): + value = self._inner.get_field(pos) + if pos in self._blob_field_indices and not self._blob_as_descriptor: + return _resolve_blob_payload(value, self._file_io) + if pos in self._descriptor_field_indices and not self._blob_as_descriptor: + return _resolve_descriptor_payload(value, self._file_io) + if pos in self._view_field_indices and self._view_resolver is not None: + self._view_resolver.preload_values([value]) + return self._view_resolver.resolve(value, self._blob_as_descriptor) + return value + + def get_blob(self, pos: int): + if not (pos in self._blob_field_indices + or pos in self._descriptor_field_indices + or pos in self._view_field_indices): + return self._inner.get_blob(pos) + value = self.get_field(pos) + if value is None: + return None + if self._blob_as_descriptor: + return Blob.from_descriptor_bytes(value, file_io=self._file_io) + return Blob.from_bytes(value, self._file_io) Review Comment: [P2] Do not heuristically parse a payload that has already been materialized. For managed/descriptor/view fields, `get_field()` has already resolved the outer reference and returns opaque payload bytes. Passing those bytes through `Blob.from_bytes()` interprets a legitimate payload that happens to be a serialized v2 descriptor or `BlobViewStruct` as a second reference. I reproduced `get_field()` returning the expected descriptor-shaped payload while `get_blob().to_data()` dereferenced its URI and returned different bytes. Java preserves the typed outer `BlobRef` and never reparses its payload. This branch should return `BlobData(value)`. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
