This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new aa98e58374 [python] Parse BlobDescriptor v1/v2 bytes without
misclassifying inline payload (#9539)
aa98e58374 is described below
commit aa98e58374358681c9871d94a3e486e9ef1992cb
Author: Wenchao Wu <[email protected]>
AuthorDate: Thu Sep 3 13:41:58 2026 +0800
[python] Parse BlobDescriptor v1/v2 bytes without misclassifying inline
payload (#9539)
---
.../pypaimon/common/options/core_options.py | 8 +
paimon-python/pypaimon/table/row/blob.py | 152 ++++++++++++--
paimon-python/pypaimon/tests/blob_table_test.py | 46 ++++
paimon-python/pypaimon/tests/blob_test.py | 233 ++++++++++++++++++++-
.../write/writer/dedicated_format_writer.py | 19 +-
5 files changed, 431 insertions(+), 27 deletions(-)
diff --git a/paimon-python/pypaimon/common/options/core_options.py
b/paimon-python/pypaimon/common/options/core_options.py
index 1ae2f3187d..6e92ceb260 100644
--- a/paimon-python/pypaimon/common/options/core_options.py
+++ b/paimon-python/pypaimon/common/options/core_options.py
@@ -1238,6 +1238,14 @@ class CoreOptions:
return val
def blob_descriptor_fields(self, default=None):
+ # Do not treat blob.stored-descriptor-fields as a layout switch.
+ # Python master ignored that key and wrote dedicated .blob payloads;
+ # a global fallback would mis-parse those files during a rolling
+ # upgrade. The cost is that Java tables which only set the fallback
+ # key store inline descriptors, and Python returns those bytes
+ # instead of fetching payload. Migrate explicitly to
+ # blob-descriptor-field (column directives already copy the legacy
+ # key onto the canonical option).
value = self.options.get(CoreOptions.BLOB_DESCRIPTOR_FIELD, default)
return CoreOptions._parse_field_set(value)
diff --git a/paimon-python/pypaimon/table/row/blob.py
b/paimon-python/pypaimon/table/row/blob.py
index 4988f61c24..fb700721de 100644
--- a/paimon-python/pypaimon/table/row/blob.py
+++ b/paimon-python/pypaimon/table/row/blob.py
@@ -28,6 +28,8 @@ from pypaimon.common.uri_reader import UriReader,
FileUriReader
class BlobDescriptor:
CURRENT_VERSION = 2
MAGIC = 0x424C4F4244455343 # "BLOBDESC"
+ # v1 wire: version (1) + uri_length (4) + offset (8) + length (8)
+ _V1_MIN_WIRE_SIZE = 1 + 4 + 16
def __init__(self, uri: str, offset: int, length: int):
self._version = self.CURRENT_VERSION
@@ -54,13 +56,13 @@ class BlobDescriptor:
def serialize(self) -> bytes:
uri_bytes = self._uri.encode('utf-8')
uri_length = len(uri_bytes)
- data = struct.pack('<B', self._version) # version (1 byte)
- if self._version > 1:
- data += struct.pack('<Q', self.MAGIC) # magic (8 bytes, unsigned)
- data += struct.pack('<I', uri_length) # uri length (4 bytes)
- data += uri_bytes # uri bytes
- data += struct.pack('<q', self._offset) # offset (8 bytes, signed)
- data += struct.pack('<q', self._length) # length (8 bytes, signed)
+ # Always write CURRENT_VERSION with magic, matching Java
BlobDescriptor.serialize().
+ data = struct.pack('<B', self.CURRENT_VERSION)
+ data += struct.pack('<Q', self.MAGIC)
+ data += struct.pack('<i', uri_length)
+ data += uri_bytes
+ data += struct.pack('<q', self._offset)
+ data += struct.pack('<q', self._length)
return data
@classmethod
@@ -100,8 +102,12 @@ class BlobDescriptor:
# Read URI length
if offset + 4 > len(data):
raise ValueError("Invalid BlobDescriptor data: too short")
- uri_length = struct.unpack('<I', data[offset:offset + 4])[0]
+ uri_length = struct.unpack('<i', data[offset:offset + 4])[0]
offset += 4
+ if uri_length < 0:
+ raise ValueError(
+ f"Invalid BlobDescriptor data: negative URI length:
{uri_length}"
+ )
# Read URI bytes
if offset + uri_length > len(data):
@@ -124,6 +130,51 @@ class BlobDescriptor:
descriptor._version = version
return descriptor
+ @classmethod
+ def parse_if_serialized(cls, data: bytes) -> Optional['BlobDescriptor']:
+ """Parse when data is exactly a serialized descriptor (no trailing
bytes).
+
+ Dispatches through :class:`BlobDescriptorSerde` so an exact
+ :class:`VideoFrameDescriptor` is accepted before the ordinary v1/v2
+ BlobDescriptor length check. Unlike :meth:`is_blob_descriptor` (v2
+ magic header only), this accepts v1 descriptors without a magic
+ prefix. Unlike ordinary :meth:`deserialize`, the encoded length must
+ match the buffer exactly. Still heuristic: arbitrary inline blob
+ bytes could theoretically match.
+ """
+ if not isinstance(data, (bytes, bytearray)):
+ return None
+ return BlobDescriptorSerde.parse_if_serialized(bytes(data))
+
+ @classmethod
+ def _parse_ordinary_if_serialized(cls, raw: bytes) ->
Optional['BlobDescriptor']:
+ if len(raw) < cls._V1_MIN_WIRE_SIZE:
+ 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('<Q', raw[offset:offset + 8])[0]
+ if magic != cls.MAGIC:
+ return None
+ offset += 8
+ if offset + 4 > len(raw):
+ return None
+ uri_length = struct.unpack('<i', raw[offset:offset + 4])[0]
+ if uri_length < 0:
+ return None
+ total = offset + 4 + uri_length + 16
+ if total != len(raw):
+ return None
+ return cls._deserialize(raw)
+ except (ValueError, struct.error, UnicodeDecodeError):
+ return None
+
@classmethod
def is_blob_descriptor(cls, data: bytes) -> bool:
if not isinstance(data, (bytes, bytearray)):
@@ -296,6 +347,19 @@ class BlobDescriptorSerde:
return VideoFrameDescriptor.deserialize(data)
return BlobDescriptor._deserialize(data)
+ @staticmethod
+ def parse_if_serialized(data: bytes) -> Optional[BlobDescriptor]:
+ """Exact-length parse for any persisted BlobDescriptor wire type."""
+ if not isinstance(data, (bytes, bytearray)):
+ return None
+ raw = bytes(data)
+ if VideoFrameDescriptor.is_video_frame_descriptor(raw):
+ try:
+ return VideoFrameDescriptor.deserialize(raw)
+ except (ValueError, struct.error, UnicodeDecodeError):
+ return None
+ return BlobDescriptor._parse_ordinary_if_serialized(raw)
+
class BlobViewStruct:
CURRENT_VERSION = 1
@@ -520,6 +584,49 @@ class Blob(ABC):
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.
+
+ Parsing uses :meth:`BlobDescriptor.deserialize`, matching Java: a
+ valid v1/v2 prefix is accepted and trailing bytes after that prefix
+ are ignored. This is not a detector; garbage that happens to look
+ like a v1 prefix can produce a BlobRef with a nonsense URI.
+ Bytes that are not a parseable prefix raise :class:`ValueError`.
+ """
+ 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)
@@ -535,20 +642,16 @@ class Blob(ABC):
data = bytes(data)
if BlobViewStruct.is_blob_view_struct(data):
return Blob.from_view(BlobViewStruct.deserialize(data))
- is_descriptor = BlobDescriptorSerde.is_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 = BlobDescriptorSerde.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 BlobDescriptorSerde.is_descriptor(data) or not allow_blob_data:
+ try:
+ descriptor = BlobDescriptor.deserialize(data)
+ except (ValueError, struct.error, UnicodeDecodeError) as exc:
+ raise ValueError(
+ "Expected BlobDescriptor bytes, got raw bytes"
+ + ("" if allow_blob_data else " (allow_blob_data=False)")
+ ) from exc
+ return Blob._blob_ref_from_descriptor(
+ descriptor, file_io=file_io,
uri_reader_factory=uri_reader_factory)
return BlobData(data)
@@ -637,6 +740,11 @@ class BlobRef(Blob):
def to_descriptor(self) -> BlobDescriptor:
return self._descriptor
+ @property
+ def uri_reader(self) -> UriReader:
+ """UriReader used to fetch this blob's payload."""
+ return self._uri_reader
+
def new_input_stream(self) -> BinaryIO:
uri = self._descriptor.uri
offset = self._descriptor.offset
diff --git a/paimon-python/pypaimon/tests/blob_table_test.py
b/paimon-python/pypaimon/tests/blob_table_test.py
index c08289b4ef..1ab02b56ea 100755
--- a/paimon-python/pypaimon/tests/blob_table_test.py
+++ b/paimon-python/pypaimon/tests/blob_table_test.py
@@ -5247,6 +5247,52 @@ class DedicatedFormatWriterTest(unittest.TestCase):
self.assertEqual(result['id'], list(range(2000)))
self.assertEqual(result['name'], ['updated'] * 2000)
+ def test_legacy_stored_descriptor_fields_keeps_dedicated_blob_layout(self):
+ """blob.stored-descriptor-fields must not switch Python to inline
descriptors.
+
+ Master ignored that key and wrote dedicated .blob payloads. Head write
+ with the same option must keep that layout so old readers still see
+ payloads, and head read must not fail-fast on those bytes.
+ """
+ from pypaimon import Schema
+
+ pa_schema = pa.schema([
+ ('id', pa.int32()),
+ ('picture', pa.large_binary()),
+ ])
+ schema = Schema.from_pyarrow_schema(
+ pa_schema,
+ options={
+ 'row-tracking.enabled': 'true',
+ 'data-evolution.enabled': 'true',
+ 'blob.stored-descriptor-fields': 'picture',
+ }
+ )
+ self.catalog.create_table(
+ 'test_db.legacy_stored_descriptor_fields', schema, False)
+ table =
self.catalog.get_table('test_db.legacy_stored_descriptor_fields')
+
+ payload = b'legacy-dedicated-blob-payload'
+ write_builder = table.new_batch_write_builder()
+ writer = write_builder.new_write()
+ writer.write_arrow(pa.Table.from_pydict({
+ 'id': [1],
+ 'picture': [payload],
+ }, schema=pa_schema))
+ commit_messages = writer.prepare_commit()
+ write_builder.new_commit().commit(commit_messages)
+ writer.close()
+
+ all_files = [f for msg in commit_messages for f in msg.new_files]
+ blob_files = [f for f in all_files if f.file_name.endswith('.blob')]
+ self.assertGreaterEqual(len(blob_files), 1)
+ self.assertTrue(all(f.write_cols == ['picture'] for f in blob_files))
+
+ result = table.new_read_builder().new_read().to_arrow(
+ table.new_read_builder().new_scan().plan().splits())
+ self.assertEqual(result.num_rows, 1)
+ self.assertEqual(result.column('picture').to_pylist()[0], payload)
+
class GetBlobTest(unittest.TestCase):
diff --git a/paimon-python/pypaimon/tests/blob_test.py
b/paimon-python/pypaimon/tests/blob_test.py
index c0ba12bdfe..456d8dc2b2 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -41,7 +41,16 @@ from pypaimon.manifest.schema.data_file_meta import
DataFileMeta
from pypaimon.read.reader.concat_batch_reader import BlobFallbackBatchReader
from pypaimon.read.reader.format_blob_reader import BlobRecordIterator,
FormatBlobReader
from pypaimon.schema.data_types import ArrayType, AtomicType, DataField,
MapType
-from pypaimon.table.row.blob import Blob, BlobData, BlobRef, BlobDescriptor,
BlobViewStruct, BlobView
+from pypaimon.table.row.blob import (
+ Blob,
+ BlobData,
+ BlobRef,
+ BlobDescriptor,
+ BlobDescriptorSerde,
+ BlobViewStruct,
+ BlobView,
+ VideoFrameDescriptor,
+)
from pypaimon.table.row.generic_row import GenericRowDeserializer,
GenericRowSerializer, GenericRow
from pypaimon.table.row.row_kind import RowKind
from pypaimon.utils.range import Range
@@ -1282,6 +1291,14 @@ class BlobTest(unittest.TestCase):
BlobDescriptor.deserialize(incomplete_data)
self.assertIn("URI length exceeds data size", str(context.exception))
+ # Java reads uri length as a signed int and rejects negatives.
+ negative_uri = bytearray(valid_descriptor.serialize())
+ struct.pack_into('<i', negative_uri, 1 + 8, -1)
+ with self.assertRaises(ValueError) as context:
+ BlobDescriptor.deserialize(bytes(negative_uri))
+ self.assertIn("negative URI length", str(context.exception))
+
self.assertIsNone(BlobDescriptor.parse_if_serialized(bytes(negative_uri)))
+
def test_blob_descriptor_equality_and_hashing(self):
"""Test BlobDescriptor equality and hashing."""
# Create identical descriptors
@@ -1411,6 +1428,16 @@ class BlobTest(unittest.TestCase):
)
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('<I', 5)
+ + b"hello"
+ + struct.pack('<q', 0)
+ + struct.pack('<q', 5)
+ + b"\xff"
+ )
v2_magic_only = bytes([2]) + struct.pack('<Q', BlobDescriptor.MAGIC)
self.assertTrue(BlobDescriptor.is_blob_descriptor(descriptor_v2.serialize()))
@@ -1421,6 +1448,210 @@ class BlobTest(unittest.TestCase):
self.assertFalse(BlobDescriptor.is_blob_descriptor(fake_v1_prefix))
self.assertFalse(BlobDescriptor.is_blob_descriptor(b"tiny"))
+
self.assertIsNotNone(BlobDescriptor.parse_if_serialized(descriptor_v1_bytes))
+
self.assertIsNotNone(BlobDescriptor.parse_if_serialized(descriptor_v2.serialize()))
+ self.assertIsNone(BlobDescriptor.parse_if_serialized(random_bytes))
+ self.assertIsNone(BlobDescriptor.parse_if_serialized(fake_v1_prefix))
+ self.assertIsNone(BlobDescriptor.parse_if_serialized(v1_shaped_inline))
+ self.assertIsNone(BlobDescriptor.parse_if_serialized(b"tiny"))
+
+ video = VideoFrameDescriptor("file:///v.mp4", 0, 10, 2)
+ video_bytes = video.serialize()
+ self.assertEqual(video_bytes,
BlobDescriptor.deserialize(video_bytes).serialize())
+ self.assertEqual(video,
BlobDescriptor.parse_if_serialized(video_bytes))
+ self.assertEqual(video,
BlobDescriptorSerde.parse_if_serialized(video_bytes))
+ self.assertIsNone(BlobDescriptor.parse_if_serialized(video_bytes +
b"x"))
+
+ 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):
+ uri = b"file:///tmp/blob.bin"
+ serialized_v1 = (
+ bytes([1])
+ + struct.pack('<I', len(uri))
+ + uri
+ + struct.pack('<q', 0)
+ + struct.pack('<q', 10)
+ )
+ padded = serialized_v1 + b"extra"
+ blob = Blob.from_descriptor_bytes(padded, file_io=LocalFileIO())
+ self.assertIsInstance(blob, BlobRef)
+ self.assertEqual(blob.to_descriptor().uri, "file:///tmp/blob.bin")
+
+ def
test_from_descriptor_bytes_accepts_v2_descriptor_with_trailing_bytes(self):
+ descriptor = BlobDescriptor("file:///tmp/blob.bin", 0, 10)
+ padded = descriptor.serialize() + b"extra"
+ blob = Blob.from_descriptor_bytes(padded, file_io=LocalFileIO())
+ self.assertIsInstance(blob, BlobRef)
+ self.assertEqual(blob.to_descriptor().uri, descriptor.uri)
+
+ def test_from_bytes_allow_blob_data_false_rejects_raw_bytes(self):
+ with self.assertRaises(ValueError) as ctx:
+ Blob.from_bytes(b"hello blob", allow_blob_data=False)
+ self.assertIn("Expected BlobDescriptor bytes", str(ctx.exception))
+ self.assertIn("allow_blob_data=False", str(ctx.exception))
+
+ def
test_from_bytes_allow_blob_data_false_requires_file_io_for_valid_descriptor(self):
+ descriptor = BlobDescriptor("file:///tmp/blob.bin", 0, 10)
+ with self.assertRaises(ValueError) as ctx:
+ Blob.from_bytes(descriptor.serialize(), allow_blob_data=False)
+ self.assertIn("file_io is required", str(ctx.exception))
+ self.assertNotIn("allow_blob_data=False", str(ctx.exception))
+
+ def test_from_bytes_allow_blob_data_false_accepts_v1_descriptor(self):
+ data = b"actual blob content"
+ with tempfile.TemporaryDirectory() as tmp_dir:
+ blob_path = os.path.join(tmp_dir, "blob.bin")
+ with open(blob_path, 'wb') as f:
+ f.write(data)
+ uri = blob_path.encode('utf-8')
+ serialized_v1 = (
+ bytes([1])
+ + struct.pack('<I', len(uri))
+ + uri
+ + struct.pack('<q', 0)
+ + struct.pack('<q', len(data))
+ )
+ file_io = FileIO.get(f"file://{tmp_dir}", {})
+ blob = Blob.from_bytes(serialized_v1, file_io=file_io,
allow_blob_data=False)
+ self.assertIsInstance(blob, BlobRef)
+ self.assertEqual(blob.to_data(), data)
+
+ def test_from_descriptor_bytes_with_v1_descriptor_bytes(self):
+ data = b"actual blob content"
+ with tempfile.TemporaryDirectory() as tmp_dir:
+ blob_path = os.path.join(tmp_dir, "blob.bin")
+ with open(blob_path, 'wb') as f:
+ f.write(data)
+ uri = blob_path.encode('utf-8')
+ serialized_v1 = (
+ bytes([1])
+ + struct.pack('<I', len(uri))
+ + uri
+ + struct.pack('<q', 0)
+ + struct.pack('<q', len(data))
+ )
+ file_io = FileIO.get(f"file://{tmp_dir}", {})
+ blob = Blob.from_descriptor_bytes(serialized_v1, file_io)
+ self.assertIsInstance(blob, BlobRef)
+ self.assertEqual(blob.to_data(), data)
+
+ def test_from_bytes_does_not_misparse_v1_shaped_inline_payload(self):
+ # Exact-length inline bytes that match the v1 descriptor wire layout
must
+ # stay as BlobData when parsed via the heuristic from_bytes entry
point.
+ v1_shaped_inline = (
+ bytes([1])
+ + struct.pack('<I', 5)
+ + b"hello"
+ + struct.pack('<q', 0)
+ + struct.pack('<q', 5)
+ )
+ blob = Blob.from_bytes(v1_shaped_inline)
+ self.assertIsInstance(blob, BlobData)
+ self.assertEqual(blob.to_data(), v1_shaped_inline)
+
+ def test_blob_descriptor_serialize_uses_current_version(self):
+ uri = b"file:///tmp/blob.bin"
+ serialized_v1 = (
+ bytes([1])
+ + struct.pack('<I', len(uri))
+ + uri
+ + struct.pack('<q', 0)
+ + struct.pack('<q', 10)
+ )
+ descriptor = BlobDescriptor.deserialize(serialized_v1)
+ self.assertEqual(descriptor.version, 1)
+
+ serialized = descriptor.serialize()
+ self.assertEqual(serialized[0], BlobDescriptor.CURRENT_VERSION)
+ self.assertTrue(BlobDescriptor.is_blob_descriptor(serialized))
+ self.assertEqual(
+ BlobDescriptor.deserialize(serialized),
+ BlobDescriptor("file:///tmp/blob.bin", 0, 10),
+ )
+
+ def test_blob_descriptor_fields_ignores_legacy_stored_key(self):
+ from pypaimon.common.options.core_options import CoreOptions
+
+ # Python master ignored this key and wrote dedicated .blob files.
+ # A global fallback would break rolling upgrades.
+ legacy_only = CoreOptions(
+ Options({"blob.stored-descriptor-fields": "legacy_col"}))
+ self.assertEqual(set(), legacy_only.blob_descriptor_fields())
+
+ canonical_wins = CoreOptions(Options({
+ "blob-descriptor-field": "canon",
+ "blob.stored-descriptor-fields": "legacy_col",
+ }))
+ self.assertEqual({"canon"}, canonical_wins.blob_descriptor_fields())
+
+ blank_canonical = CoreOptions(Options({
+ "blob-descriptor-field": "",
+ "blob.stored-descriptor-fields": "legacy_col",
+ }))
+ self.assertEqual(set(), blank_canonical.blob_descriptor_fields())
+
+ def test_dedicated_writer_accepts_exact_v1_descriptor_bytes(self):
+ from pypaimon.write.writer.dedicated_format_writer import (
+ DedicatedFormatWriter)
+
+ uri = b"file:///tmp/blob.bin"
+ v1 = (
+ bytes([1])
+ + struct.pack('<I', len(uri))
+ + uri
+ + struct.pack('<q', 0)
+ + struct.pack('<q', 10)
+ )
+ writer = object.__new__(DedicatedFormatWriter)
+ writer.blob_inline_fields = {'payload'}
+ writer.blob_descriptor_fields = {'payload'}
+ writer.blob_view_fields = set()
+ v2 = BlobDescriptor("file:///tmp/blob.bin", 0, 10).serialize()
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([v2], type=pa.large_binary())], names=['payload']))
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([v1], type=pa.large_binary())], names=['payload']))
+
+ padded = pa.RecordBatch.from_arrays(
+ [pa.array([v1 + b"x"], type=pa.large_binary())], names=['payload'])
+ with self.assertRaisesRegex(ValueError, "trailing bytes"):
+ writer._validate_inline_stored_fields_input(padded)
+
+ v0 = bytes([0]) + v1[1:]
+ with self.assertRaisesRegex(ValueError, r"in \[1, 2\], but found 0"):
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([v0], type=pa.large_binary())],
names=['payload']))
+
+ v3 = bytes([3]) + v2[1:]
+ with self.assertRaisesRegex(ValueError, r"in \[1, 2\], but found 3"):
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([v3], type=pa.large_binary())],
names=['payload']))
+
+ # write_row() keeps raw bytes from _normal_row_value(); the same
+ # validator must accept exact v1 and reject padding.
+ self.assertEqual(writer._normal_row_value("payload", v1), v1)
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([v1], type=pa.large_binary())], names=["payload"]))
+
+ video_bytes = VideoFrameDescriptor("file:///v.mp4", 0, 10,
2).serialize()
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([video_bytes], type=pa.large_binary())],
names=["payload"]))
+ with self.assertRaisesRegex(ValueError, "serialized"):
+ writer._validate_inline_stored_fields_input(
+ pa.RecordBatch.from_arrays(
+ [pa.array([video_bytes + b"x"], type=pa.large_binary())],
+ names=["payload"]))
+
class BlobEndToEndTest(unittest.TestCase):
"""End-to-end tests for blob functionality with schema definition, file
writing, and reading."""
diff --git a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
index 8b2097741f..9c80ec281a 100644
--- a/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
+++ b/paimon-python/pypaimon/write/writer/dedicated_format_writer.py
@@ -296,6 +296,7 @@ class DedicatedFormatWriter(DataWriter):
self.table.table_schema.fields,
self.normal_column_names,
).to_batches()[0]
+ self._validate_inline_stored_fields_input(normal_data)
processed_normal = self._process_normal_data(normal_data)
if processed_normal is not None:
self._normal_buffer.append(processed_normal)
@@ -436,16 +437,26 @@ class DedicatedFormatWriter(DataWriter):
"blob-descriptor-field requires blob field value to be
a serialized "
"BlobDescriptor."
)
+ descriptor_bytes = bytes(value)
+ if descriptor_bytes:
+ version = descriptor_bytes[0]
+ if version < 1 or version > BlobDescriptor.CURRENT_VERSION:
+ raise ValueError(
+ f"blob-descriptor-field requires BlobDescriptor
version "
+ f"in [1, {BlobDescriptor.CURRENT_VERSION}], but
found "
+ f"{version}."
+ )
try:
- descriptor_bytes = bytes(value)
- descriptor = BlobDescriptor.deserialize(descriptor_bytes)
- if descriptor.serialize() != descriptor_bytes:
- raise ValueError("Descriptor payload contains trailing
bytes.")
+ BlobDescriptor.deserialize(descriptor_bytes)
except Exception as e:
raise ValueError(
"blob-descriptor-field requires blob field value to be
a serialized "
"BlobDescriptor."
) from e
+ # serialize() always emits CURRENT_VERSION, so a round-trip
+ # would reject exact v1 bytes. Check exact wire length instead.
+ if BlobDescriptor.parse_if_serialized(descriptor_bytes) is
None:
+ raise ValueError("Descriptor payload contains trailing
bytes.")
for field_name in self.blob_view_fields:
if field_name not in data.schema.names: