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 57201293f7 [python] Parallelize fallback BLOB reads (#8620)
57201293f7 is described below
commit 57201293f7559b05c25250c2346d433563323aa7
Author: umi <[email protected]>
AuthorDate: Thu Jul 16 13:39:57 2026 +0800
[python] Parallelize fallback BLOB reads (#8620)
---
.../pypaimon/read/reader/concat_batch_reader.py | 62 ++++-
paimon-python/pypaimon/read/split_read.py | 7 +-
paimon-python/pypaimon/tests/blob_test.py | 269 ++++++++++++++++++++-
3 files changed, 331 insertions(+), 7 deletions(-)
diff --git a/paimon-python/pypaimon/read/reader/concat_batch_reader.py
b/paimon-python/pypaimon/read/reader/concat_batch_reader.py
index f018ccf756..46b9c6df0c 100644
--- a/paimon-python/pypaimon/read/reader/concat_batch_reader.py
+++ b/paimon-python/pypaimon/read/reader/concat_batch_reader.py
@@ -260,7 +260,8 @@ class BlobFallbackBatchReader(RecordBatchReader):
def __init__(self, file_reader_suppliers: List[Tuple[DataFileMeta,
Callable]],
field_name: str, output_type, row_ranges:
Optional[List[Range]] = None,
- blob_as_descriptor: bool = False, deletion_vector=None,
batch_size: int = 1024):
+ blob_as_descriptor: bool = False, deletion_vector=None,
batch_size: int = 1024,
+ blob_parallelism: int = 1):
self._file_reader_suppliers = file_reader_suppliers
self._field_name = field_name
self._output_type = output_type
@@ -279,6 +280,8 @@ class BlobFallbackBatchReader(RecordBatchReader):
self._deletion_vector_range, self._deletion_vector =
deletion_vector
self._returned = False
self._batch_size = max(1, batch_size)
+ self._blob_parallelism = max(1, blob_parallelism)
+ self._file_io = None
self._target_ranges = self._compute_target_ranges()
self._target_range_index = 0
self._next_row_id = (
@@ -297,6 +300,9 @@ class BlobFallbackBatchReader(RecordBatchReader):
if not batch_row_ids:
return None
+ resolve_blobs_concurrently = (
+ self._blob_parallelism > 1 and not self._blob_as_descriptor
+ )
groups: Dict[int, Dict[int, Tuple[object, bool]]] = {}
batch_first = batch_row_ids[0]
@@ -308,6 +314,8 @@ class BlobFallbackBatchReader(RecordBatchReader):
if not blob_values:
continue
group = groups.setdefault(state.file.max_sequence_number, {})
+ if resolve_blobs_concurrently and self._file_io is None:
+ self._file_io = state.reader._file_io
for row_id, blob in blob_values.items():
if row_id in group:
raise ValueError(
@@ -318,10 +326,18 @@ class BlobFallbackBatchReader(RecordBatchReader):
elif blob is Blob.PLACE_HOLDER or blob is
Blob.ARRAY_PLACE_HOLDER:
group[row_id] = (None, True)
elif self._is_array_blob:
- group[row_id] = (self._array_value_for_arrow(blob), False)
+ if resolve_blobs_concurrently:
+ group[row_id] = (blob, False)
+ else:
+ group[row_id] = (self._array_value_for_arrow(blob),
False)
else:
if self._blob_as_descriptor:
group[row_id] = (blob.to_descriptor().serialize(),
False)
+ elif resolve_blobs_concurrently:
+ # Keep values lazy until fallback selects the newest
+ # non-placeholder version for each row. Otherwise older
+ # overridden BLOB versions would be read unnecessarily.
+ group[row_id] = (blob, False)
else:
group[row_id] = (blob.to_data(), False)
if state.selected_range_index >= len(state.selected_ranges):
@@ -345,6 +361,9 @@ class BlobFallbackBatchReader(RecordBatchReader):
if not found:
raise ValueError("All blob files at the same row id store a
placeholder.")
+ if resolve_blobs_concurrently:
+ result = self._resolve_selected_blobs(result)
+
return pa.RecordBatch.from_arrays(
[pa.array(result, type=self._output_type)],
names=[self._field_name],
@@ -361,6 +380,41 @@ class BlobFallbackBatchReader(RecordBatchReader):
result.append(blob.to_data())
return result
+ def _resolve_selected_blobs(self, values: List[object]) -> List[object]:
+ """Materialize selected scalar or array BLOBs with the shared
FileIO."""
+ resolved = []
+ indexed_blobs: List[Tuple[int, Optional[int], Blob]] = []
+ for row_index, value in enumerate(values):
+ if self._is_array_blob:
+ if value is None:
+ resolved.append(None)
+ continue
+ array_values = []
+ for element_index, blob in enumerate(value):
+ if blob is None:
+ array_values.append(None)
+ else:
+ array_values.append(None)
+ indexed_blobs.append((row_index, element_index, blob))
+ resolved.append(array_values)
+ elif isinstance(value, Blob):
+ resolved.append(None)
+ indexed_blobs.append((row_index, None, value))
+ else:
+ resolved.append(value)
+
+ if not indexed_blobs:
+ return resolved
+
+ bodies = self._file_io.read_blobs_concurrent(
+ [blob for _, _, blob in indexed_blobs], self._blob_parallelism)
+ for (row_index, element_index, _), body in zip(indexed_blobs, bodies):
+ if element_index is None:
+ resolved[row_index] = body
+ else:
+ resolved[row_index][element_index] = body
+ return resolved
+
def _compute_target_ranges(self) -> List[Range]:
ranges = Range.sort_and_merge_overlap([
file.row_id_range()
@@ -459,7 +513,9 @@ class BlobFallbackBatchReader(RecordBatchReader):
blob_offsets,
self._data_field,
reader._input_stream,
- blob_as_descriptor=self._blob_as_descriptor,
+ blob_as_descriptor=(
+ self._blob_as_descriptor or self._blob_parallelism > 1
+ ),
)
blobs = []
diff --git a/paimon-python/pypaimon/read/split_read.py
b/paimon-python/pypaimon/read/split_read.py
index 0b01ae4507..ed690a78da 100644
--- a/paimon-python/pypaimon/read/split_read.py
+++ b/paimon-python/pypaimon/read/split_read.py
@@ -284,7 +284,7 @@ class SplitRead(ABC):
raise NotImplementedError(
"Nested-field projection is not supported on BLOB files")
blob_as_descriptor =
CoreOptions.blob_as_descriptor(self.table.options)
- blob_parallelism = getattr(self, '_blob_parallelism', 1)
+ blob_parallelism = self._blob_parallelism
format_reader = FormatBlobReader(self.table.file_io, file_path,
read_file_fields,
self.read_fields,
read_arrow_predicate, blob_as_descriptor,
batch_size=batch_size,
@@ -1005,7 +1005,7 @@ class DataEvolutionSplitRead(SplitRead):
self.table.options))
or (not CoreOptions.blob_as_descriptor(self.table.options)
and
CoreOptions.blob_descriptor_fields(self.table.options))):
- blob_parallelism = getattr(self, '_blob_parallelism', 1)
+ blob_parallelism = self._blob_parallelism
reader = BlobInlineConvertReader(
reader, self.table,
prescan_reader_factory=lambda names:
self._create_prescan_reader(names),
@@ -1284,6 +1284,7 @@ class DataEvolutionSplitRead(SplitRead):
CoreOptions.blob_as_descriptor(self.table.options),
deletion_vector=deletion_vector,
batch_size=batch_size,
+ blob_parallelism=self._blob_parallelism,
)
else:
# Create concatenated reader for multiple files
@@ -1328,7 +1329,7 @@ class DataEvolutionSplitRead(SplitRead):
return None
file_path = file.external_path if file.external_path else
file.file_path
- blob_parallelism = getattr(self, '_blob_parallelism', 1)
+ blob_parallelism = self._blob_parallelism
return FormatBlobReader(
self.table.file_io,
file_path,
diff --git a/paimon-python/pypaimon/tests/blob_test.py
b/paimon-python/pypaimon/tests/blob_test.py
index 68b72da696..a743442b5b 100644
--- a/paimon-python/pypaimon/tests/blob_test.py
+++ b/paimon-python/pypaimon/tests/blob_test.py
@@ -22,6 +22,7 @@ import struct
import tempfile
import unittest
from pathlib import Path
+from unittest.mock import patch
import pyarrow as pa
@@ -289,6 +290,12 @@ class BlobTest(unittest.TestCase):
def test_blob_fallback_batch_reader_respects_batch_size(self):
created_readers = []
+ class DescriptorBlobFallbackBatchReader(BlobFallbackBatchReader):
+ def _resolve_selected_blobs(self, values):
+ raise AssertionError(
+ "Descriptor reads should not materialize BLOB data."
+ )
+
class FakeBlobReader:
def __init__(self):
self._file_io = None
@@ -322,12 +329,13 @@ class BlobTest(unittest.TestCase):
first_row_id=10,
file_path="fake.blob",
)
- reader = BlobFallbackBatchReader(
+ reader = DescriptorBlobFallbackBatchReader(
[(data_file, supplier)],
"picture",
pa.large_binary(),
blob_as_descriptor=True,
batch_size=2,
+ blob_parallelism=4,
)
first = reader.read_arrow_batch()
@@ -531,6 +539,198 @@ class BlobTest(unittest.TestCase):
reader.close()
self.assertTrue(created_by_file["third.blob"][0].closed)
+ def
test_blob_fallback_batch_reader_materializes_selected_values_in_parallel(self):
+ class RecordingFileIO:
+ def __init__(self):
+ self.calls = []
+
+ def read_blobs_concurrent(self, blobs, parallelism):
+ descriptors = [blob.to_descriptor() for blob in blobs]
+ self.calls.append((descriptors, parallelism))
+ return [
+ "{}:{}".format(descriptor.uri, descriptor.offset).encode()
+ for descriptor in descriptors
+ ]
+
+ class FakeBlobReader:
+ def __init__(self, file_io, file_path, blob_lengths, blob_offsets):
+ self._file_io = file_io
+ self.file_path = file_path
+ self.blob_lengths = blob_lengths
+ self.blob_offsets = blob_offsets
+ self._input_stream = None
+
+ def close(self):
+ pass
+
+ def data_file(name, max_sequence_number):
+ return DataFileMeta(
+ file_name=name,
+ file_size=0,
+ row_count=3,
+ min_key=None,
+ max_key=None,
+ key_stats=None,
+ value_stats=None,
+ min_sequence_number=max_sequence_number,
+ max_sequence_number=max_sequence_number,
+ schema_id=0,
+ level=0,
+ extra_files=[],
+ first_row_id=0,
+ file_path=name,
+ )
+
+ file_io = RecordingFileIO()
+ old_file = data_file("old.blob", 1)
+ new_file = data_file("new.blob", 2)
+ reader = BlobFallbackBatchReader(
+ [
+ (
+ old_file,
+ lambda: FakeBlobReader(
+ file_io, "old.blob", [20, 20, 20], [0, 100, 200]
+ ),
+ ),
+ (
+ new_file,
+ lambda: FakeBlobReader(
+ file_io, "new.blob", [-2, 20, -2], [-1, 1000, -1]
+ ),
+ ),
+ ],
+ "picture",
+ pa.large_binary(),
+ batch_size=3,
+ blob_parallelism=4,
+ )
+
+ batch = reader.read_arrow_batch()
+
+ self.assertEqual(
+ [b"old.blob:4", b"new.blob:1004", b"old.blob:204"],
+ batch.column("picture").to_pylist(),
+ )
+ self.assertEqual(1, len(file_io.calls))
+ descriptors, parallelism = file_io.calls[0]
+ self.assertEqual(4, parallelism)
+ self.assertEqual(
+ [("old.blob", 4), ("new.blob", 1004), ("old.blob", 204)],
+ [(descriptor.uri, descriptor.offset) for descriptor in
descriptors],
+ )
+
+ def
test_blob_fallback_batch_reader_materializes_selected_array_values_in_parallel(self):
+ from pypaimon.write.blob_format_writer import BlobFormatWriter
+
+ class RecordingFileIO(LocalFileIO):
+ def __init__(self, path, options):
+ super().__init__(path, options)
+ self.calls = []
+
+ def read_blobs_concurrent(self, blobs, parallelism):
+ self.calls.append((
+ [blob.to_descriptor() for blob in blobs],
+ parallelism,
+ ))
+ return super().read_blobs_concurrent(blobs, parallelism)
+
+ field = DataField(
+ 0,
+ "pictures",
+ ArrayType(True, AtomicType("BLOB")),
+ )
+ file_io = RecordingFileIO(self.temp_dir, Options({}))
+
+ def write_blob_file(name, values):
+ path = os.path.join(self.temp_dir, name)
+ with open(path, 'wb') as output:
+ writer = BlobFormatWriter(output)
+ for value in values:
+ writer.add_element(
+ GenericRow([value], [field], RowKind.INSERT)
+ )
+ writer.close()
+ return path
+
+ old_path = write_blob_file(
+ "old-array.blob",
+ [
+ [BlobData(b"old-0"), None],
+ [BlobData(b"old-1")],
+ [BlobData(b"old-2a"), BlobData(b"old-2b")],
+ ],
+ )
+ new_path = write_blob_file(
+ "new-array.blob",
+ [
+ Blob.ARRAY_PLACE_HOLDER,
+ [BlobData(b"new-1a"), None, BlobData(b"new-1b")],
+ Blob.ARRAY_PLACE_HOLDER,
+ ],
+ )
+
+ def data_file(path, max_sequence_number):
+ return DataFileMeta(
+ file_name=os.path.basename(path),
+ file_size=os.path.getsize(path),
+ row_count=3,
+ min_key=None,
+ max_key=None,
+ key_stats=None,
+ value_stats=None,
+ min_sequence_number=max_sequence_number,
+ max_sequence_number=max_sequence_number,
+ schema_id=0,
+ level=0,
+ extra_files=[],
+ first_row_id=0,
+ file_path=path,
+ )
+
+ def supplier(path):
+ return lambda: FormatBlobReader(
+ file_io=file_io,
+ file_path=path,
+ read_fields=[field.name],
+ full_fields=[field],
+ push_down_predicate=None,
+ blob_as_descriptor=False,
+ blob_parallelism=4,
+ )
+
+ reader = BlobFallbackBatchReader(
+ [
+ (data_file(old_path, 1), supplier(old_path)),
+ (data_file(new_path, 2), supplier(new_path)),
+ ],
+ field.name,
+ pa.list_(pa.large_binary()),
+ batch_size=3,
+ blob_parallelism=4,
+ )
+ try:
+ batch = reader.read_arrow_batch()
+ self.assertEqual(
+ [
+ [b"old-0", None],
+ [b"new-1a", None, b"new-1b"],
+ [b"old-2a", b"old-2b"],
+ ],
+ batch.column(field.name).to_pylist(),
+ )
+ self.assertIsNone(reader.read_arrow_batch())
+ finally:
+ reader.close()
+
+ self.assertEqual(1, len(file_io.calls))
+ descriptors, parallelism = file_io.calls[0]
+ self.assertEqual(4, parallelism)
+ self.assertEqual(5, len(descriptors))
+ self.assertEqual(
+ [old_path, new_path, new_path, old_path, old_path],
+ [descriptor.uri for descriptor in descriptors],
+ )
+
def test_blob_data_interface_compliance(self):
"""Test that BlobData properly implements Blob interface."""
test_data = b"interface test data"
@@ -2248,6 +2448,73 @@ class BlobParallelismTest(unittest.TestCase):
for i in range(20):
self.assertEqual(got[i], self.payloads[i])
+ def test_blob_fallback_parallelism_end_to_end(self):
+ t = self.catalog.get_table('default.bp_test')
+
+ row_id_builder = t.new_read_builder().with_projection(['id',
'_ROW_ID'])
+ row_id_result = row_id_builder.new_read().to_arrow(
+ row_id_builder.new_scan().plan().splits())
+ row_ids_by_id = dict(zip(
+ row_id_result['id'].to_pylist(),
+ row_id_result['_ROW_ID'].to_pylist(),
+ ))
+
+ updated_payload = os.urandom(512)
+ update_builder = t.new_batch_write_builder()
+ table_update = update_builder.new_update().with_update_type(['img'])
+ update_messages = table_update.update_by_arrow_with_row_id(
+ pa.Table.from_pydict({
+ '_ROW_ID': pa.array([row_ids_by_id[1]], type=pa.int64()),
+ 'img': pa.array([updated_payload], type=pa.large_binary()),
+ }))
+ update_builder.new_commit().commit(update_messages)
+
+ update_blob_files = [
+ file
+ for message in update_messages
+ for file in message.new_files
+ if file.file_name.endswith('.blob')
+ ]
+ self.assertEqual(1, len(update_blob_files))
+ blob_reader = FormatBlobReader(
+ t.file_io,
+ update_blob_files[0].file_path,
+ ['img'],
+ t.fields,
+ None,
+ False,
+ )
+ try:
+ self.assertIn(
+ FormatBlobReader.PLACE_HOLDER_LENGTH,
+ blob_reader.blob_lengths,
+ )
+ finally:
+ blob_reader.close()
+
+ rb = t.new_read_builder().with_projection(['id', 'img'])
+ splits = rb.new_scan().plan().splits()
+ serial = rb.new_read().to_arrow(splits, blob_parallelism=1)
+
+ resolve_calls = []
+ original_resolve = BlobFallbackBatchReader._resolve_selected_blobs
+
+ def tracking_resolve(reader, values):
+ resolve_calls.append(len(values))
+ return original_resolve(reader, values)
+
+ with patch.object(
+ BlobFallbackBatchReader,
+ '_resolve_selected_blobs',
+ tracking_resolve,
+ ):
+ parallel = rb.new_read().to_arrow(splits, blob_parallelism=4)
+
+ self.assertEqual(20, serial.num_rows)
+ self.assertEqual(serial.to_pydict(), parallel.to_pydict())
+ self.assertEqual(updated_payload, parallel['img'][1].as_py())
+ self.assertGreater(len(resolve_calls), 0)
+
class CapBlobParallelismTest(unittest.TestCase):
"""Peak blob threads on the parallel path (workers * blob_parallelism)