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)

Reply via email to