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 c62e049778 [python] Support ARRAY BLOB in multimodal payload reads 
(#10043)
c62e049778 is described below

commit c62e049778495c6ff552ea91eac81498afb8e79c
Author: chaoyang <[email protected]>
AuthorDate: Mon Sep 21 15:04:33 2026 +0800

    [python] Support ARRAY BLOB in multimodal payload reads (#10043)
---
 docs/docs/pypaimon/multimodal-reading.md           | 11 ++--
 paimon-python/pypaimon/multimodal/blob_read.py     | 23 ++++---
 paimon-python/pypaimon/multimodal/query.py         | 30 +++++----
 .../pypaimon/tests/multimodal_table_test.py        | 76 ++++++++++++++++++++++
 4 files changed, 116 insertions(+), 24 deletions(-)

diff --git a/docs/docs/pypaimon/multimodal-reading.md 
b/docs/docs/pypaimon/multimodal-reading.md
index f9c3730968..eb0a60239e 100644
--- a/docs/docs/pypaimon/multimodal-reading.md
+++ b/docs/docs/pypaimon/multimodal-reading.md
@@ -142,8 +142,10 @@ steps_with_imu = aligned.join_window(
 rows using concurrent, same-file coalesced ranged reads. This is much faster 
than
 a per-row loop and avoids the slow row-by-row blob resolution on data-evolution
 tables. It returns `(scalar_table, {column: rows})`, row-aligned. Scalar BLOB
-rows are `bytes | None`; `MAP<K, BLOB>` rows are `None` or ordered lists of
-`(key, bytes | None)` pairs. The scalar table drops the readable BLOB columns.
+rows are `bytes | None`; `ARRAY<BLOB>` rows are `None` or ordered lists of
+`bytes | None`; `MAP<K, BLOB>` rows are `None` or ordered lists of
+`(key, bytes | None)` pairs. Empty arrays and maps remain empty lists. The 
scalar
+table drops the readable BLOB columns, including unrequested ARRAY and MAP 
BLOBs.
 
 ```python
 scalar, blobs = (
@@ -165,12 +167,13 @@ scalar, blobs = docs.scan().where(f"id IN 
({in_clause})").read_blobs("image")
 `where()` is a SQL string, so build the `IN (...)` clause only from trusted,
 already-escaped ids -- do not interpolate untrusted external input.
 
-Read several BLOB columns at once, including `MAP<K, BLOB>`:
+Read several BLOB columns at once, including `ARRAY<BLOB>` and `MAP<K, BLOB>`:
 
 ```python
 scalar, blobs = docs.scan().where("category = 'lake'").read_blobs(
-    ["image", "audio", "renditions"]
+    ["image", "audio", "pages", "renditions"]
 )
+pages = blobs["pages"]  # row-aligned lists of page payloads, or None
 renditions = [None if row is None else dict(row) for row in 
blobs["renditions"]]
 ```
 
diff --git a/paimon-python/pypaimon/multimodal/blob_read.py 
b/paimon-python/pypaimon/multimodal/blob_read.py
index e15a3b9b94..1002a08f2f 100644
--- a/paimon-python/pypaimon/multimodal/blob_read.py
+++ b/paimon-python/pypaimon/multimodal/blob_read.py
@@ -19,13 +19,13 @@
 
 
 def fetch_blob_bodies(
-        file_io, data, blob_cols, parallelism, map_blob_cols=()):
-    """Fetch scalar and MAP BLOB payload bytes.
+        file_io, data, blob_cols, parallelism, map_blob_cols=(), 
array_blob_cols=()):
+    """Fetch scalar, MAP and ARRAY BLOB payload bytes.
 
     ``data`` is a ``dict`` mapping each BLOB column name to row-aligned cells.
     A cell may be serialized ``BlobDescriptor`` bytes, inline payload bytes,
-    ``None``, or a MAP represented by key-value pairs. Returned values preserve
-    row and MAP entry order and are grouped per column.
+    ``None``, an ARRAY of such values, or a MAP represented by key-value pairs.
+    Returned values preserve row and element order and are grouped per column.
     """
     from pypaimon.table.row.blob import (
         BlobDescriptorSerde,
@@ -38,6 +38,7 @@ def fetch_blob_bodies(
     bodies = {col: [] for col in blob_cols}
     scalar_offsets = {}
     map_blob_cols = set(map_blob_cols)
+    array_blob_cols = set(array_blob_cols)
 
     def queue_blob_fetch(value):
         index = len(ranges)
@@ -60,7 +61,8 @@ def fetch_blob_bodies(
         return index
 
     for col in blob_cols:
-        if col not in map_blob_cols:
+        is_map = col in map_blob_cols
+        if not is_map and col not in array_blob_cols:
             start = len(ranges)
             for value in data[col]:
                 queue_blob_fetch(value)
@@ -72,13 +74,13 @@ def fetch_blob_bodies(
                 bodies[col].append(None)
                 continue
 
-            entries = _map_entries(value)
+            entries = _map_entries(value) if is_map else enumerate(value)
             row_index = len(bodies[col])
             row = []
             bodies[col].append(row)
             for key, item in entries:
                 entry_index = len(row)
-                row.append((key, None))
+                row.append((key, None) if is_map else None)
                 range_index = queue_blob_fetch(item)
                 targets.append((col, row_index, entry_index, range_index))
 
@@ -93,8 +95,11 @@ def fetch_blob_bodies(
     for col, (start, end) in scalar_offsets.items():
         bodies[col] = fetched[start:end]
     for col, row_index, entry_index, index in targets:
-        key = bodies[col][row_index][entry_index][0]
-        bodies[col][row_index][entry_index] = (key, fetched[index])
+        value = fetched[index]
+        if col in map_blob_cols:
+            key = bodies[col][row_index][entry_index][0]
+            value = (key, value)
+        bodies[col][row_index][entry_index] = value
     return bodies
 
 
diff --git a/paimon-python/pypaimon/multimodal/query.py 
b/paimon-python/pypaimon/multimodal/query.py
index cb218fdc1e..fdcd33c3a0 100644
--- a/paimon-python/pypaimon/multimodal/query.py
+++ b/paimon-python/pypaimon/multimodal/query.py
@@ -22,7 +22,7 @@ import pyarrow as pa
 
 from pypaimon.common.where_parser import parse_where_clause
 from pypaimon.multimodal.blob_read import fetch_blob_bodies
-from pypaimon.schema.data_types import is_blob_type, is_map_blob_type
+from pypaimon.schema.data_types import is_array_blob_type, is_blob_type, 
is_map_blob_type
 from pypaimon.table.special_fields import SpecialFields
 
 
@@ -281,13 +281,14 @@ class ScanQuery:
     def read_blobs(
             self, columns=None, *, parallelism: int = 64
     ) -> Tuple[pa.Table, Dict[str, List[Any]]]:
-        """Materialise BLOB or MAP BLOB column(s) for the filtered rows with 
concurrent,
+        """Materialise BLOB, MAP BLOB or ARRAY BLOB columns with concurrent,
         coalesced ranged reads. Reads via blob-as-descriptor to skip the slow
         row-by-row blob resolution on multi-group data-evolution splits.
 
         ``columns`` picks the BLOB column(s) (default: all, intersected with
-        ``select(...)``). Scalar BLOB values are ``bytes|None``; MAP BLOB rows
-        are ``None`` or key-value pairs with ``bytes|None`` values. Returns a
+        ``select(...)``). Scalar BLOB values are ``bytes|None``; ARRAY BLOB 
rows
+        are ``None`` or lists of these values; MAP BLOB rows are ``None`` or
+        key-value pairs with ``bytes|None`` values. Returns a
         row-aligned ``(scalar_arrow_table, blobs_by_column)`` tuple. Use
         :meth:`stream_blobs` for a memory-bounded read.
 
@@ -297,10 +298,10 @@ class ScanQuery:
         read_builder, file_io = self._blob_descriptor_read_builder(blob_cols)
         arrow = read_builder.new_read().to_arrow(
             read_builder.new_scan().plan().splits())
-        map_blob_cols = set(blob_cols) - set(self._all_blob_columns())
+        map_blob_cols, array_blob_cols = self._nested_blob_columns()
         bodies = self._fetch_bodies(
             file_io, arrow.select(blob_cols).to_pydict(), blob_cols,
-            parallelism, map_blob_cols)
+            parallelism, map_blob_cols, array_blob_cols)
         scalar = arrow.select(self._scalar_columns(arrow.column_names))
         return scalar, bodies
 
@@ -316,12 +317,12 @@ class ScanQuery:
     def _iter_blobs(self, read_builder, file_io, blob_cols, parallelism):
         reader = read_builder.new_read().to_arrow_batch_reader(
             read_builder.new_scan().plan().splits())
-        map_blob_cols = set(blob_cols) - set(self._all_blob_columns())
+        map_blob_cols, array_blob_cols = self._nested_blob_columns()
         try:
             for batch in reader:
                 bodies = self._fetch_bodies(
                     file_io, batch.select(blob_cols).to_pydict(), blob_cols,
-                    parallelism, map_blob_cols)
+                    parallelism, map_blob_cols, array_blob_cols)
                 scalar = batch.select(self._scalar_columns(batch.schema.names))
                 yield scalar, bodies
         finally:
@@ -378,7 +379,7 @@ class ScanQuery:
 
     @staticmethod
     def _fetch_bodies(
-            file_io, data, blob_cols, parallelism, map_blob_cols=()):
+            file_io, data, blob_cols, parallelism, map_blob_cols=(), 
array_blob_cols=()):
         # Decode each descriptor to a (uri, offset, length) range and read 
them all in
         # one coalesced pass on ``file_io`` -- the read table's FileIO, which 
already
         # carries the merged DLF/OSS token. Going through Blob.from_bytes here 
would
@@ -387,7 +388,13 @@ class ScanQuery:
         # failing with "endpoint should be non-empty" / "Init credential 
failed" unless
         # the caller also passes fs.oss.* -- which users should not have to.
         return fetch_blob_bodies(
-            file_io, data, blob_cols, parallelism, map_blob_cols)
+            file_io, data, blob_cols, parallelism, map_blob_cols, 
array_blob_cols)
+
+    def _nested_blob_columns(self):
+        return (
+            [field.name for field in self._table.fields if 
is_map_blob_type(field.type)],
+            [field.name for field in self._table.fields if 
is_array_blob_type(field.type)],
+        )
 
     def _all_blob_columns(self) -> List[str]:
         return [
@@ -398,7 +405,8 @@ class ScanQuery:
     def _readable_blob_columns(self) -> List[str]:
         return [
             field.name for field in self._table.fields
-            if is_blob_type(field.type) or is_map_blob_type(field.type)
+            if (is_blob_type(field.type) or is_map_blob_type(field.type)
+                or is_array_blob_type(field.type))
         ]
 
     def _resolve_blob_columns(self, columns) -> List[str]:
diff --git a/paimon-python/pypaimon/tests/multimodal_table_test.py 
b/paimon-python/pypaimon/tests/multimodal_table_test.py
index 51edc266dc..4d7b2c2b33 100644
--- a/paimon-python/pypaimon/tests/multimodal_table_test.py
+++ b/paimon-python/pypaimon/tests/multimodal_table_test.py
@@ -1465,6 +1465,55 @@ class MultimodalTableTest(unittest.TestCase):
         _, selected = obs.scan().select(["id", "assets"]).read_blobs()
         self.assertEqual({"assets"}, set(selected))
 
+    def test_scan_read_and_stream_array_blobs(self):
+        schema = _schema({
+            "id": pa.int32(),
+            "preview": pa.large_binary(),
+            "pages": pa.list_(pa.large_binary()),
+            "assets": pa.map_(pa.string(), pa.large_binary()),
+        })
+        obs = self.conn.create_table(
+            "array_blobs", schema=schema,
+            options=dict(_PARQUET_OPTIONS, **{"read.batch-size": "2"}))
+        data = {
+            "id": [1, 2, 3, 4],
+            "preview": [b"p1", None, b"p3", b"p4"],
+            "pages": [[b"first", None, b"", b"last", b"first"], None, [], 
[b"fourth"]],
+            "assets": [[("cover", b"c1")], None, [], [("missing", None)]],
+        }
+        obs.add(pa.Table.from_pydict(data, schema=schema))
+        expected = dict(zip(data["id"], data["pages"]))
+
+        scalar, blobs = obs.scan().read_blobs()
+        self.assertEqual(["id"], scalar.column_names)
+        self.assertEqual({"preview", "pages", "assets"}, set(blobs))
+        for name in blobs:
+            self.assertEqual(dict(zip(data["id"], data[name])),
+                             dict(zip(scalar.column("id").to_pylist(), 
blobs[name])))
+
+        streamed = {}
+        for scalar_batch, blob_batch in obs.scan().stream_blobs("pages"):
+            self.assertEqual(["id"], scalar_batch.schema.names)
+            self.assertLessEqual(scalar_batch.num_rows, 2)
+            streamed.update(zip(scalar_batch.column("id").to_pylist(), 
blob_batch["pages"]))
+        self.assertEqual(expected, streamed)
+
+        # Projection hides the filter column and other BLOBs, but keeps row 
IDs.
+        query = obs.scan().select(["pages"]).with_row_id().where("id = 
4").limit(1)
+        for scalar, blobs in [query.read_blobs(), *list(query.stream_blobs())]:
+            self.assertEqual(["_ROW_ID"], scalar.schema.names)
+            self.assertEqual(1, scalar.num_rows)
+            self.assertEqual({"pages": [[b"fourth"]]}, blobs)
+
+        scalar, blobs = obs.scan().where("id < 0").read_blobs("pages")
+        self.assertEqual(0, scalar.num_rows)
+        self.assertEqual({"pages": []}, blobs)
+        self.assertEqual([], list(obs.scan().where("id < 
0").stream_blobs("pages")))
+        _, blobs = obs.scan().read_blobs(["pages", "pages"])
+        self.assertEqual({"pages"}, set(blobs))
+        scalar, _ = obs.scan().read_blobs("preview")
+        self.assertEqual(["id"], scalar.column_names)
+
     def test_scan_read_blobs_filter_column_not_selected(self):
         # The row filter must apply even when its column is not in select().
         obs = self.conn.create_table(
@@ -1826,6 +1875,33 @@ class MultimodalTableTest(unittest.TestCase):
         bodies = ScanQuery._fetch_bodies(_FakeIO(), {"img": cells}, ["img"], 8)
         self.assertEqual([b"BODY:oss://bucket/x", b"inline-blob", None], 
bodies["img"])
 
+    def test_fetch_bodies_coalesces_mixed_array_map_and_scalar_blobs(self):
+        from pypaimon.multimodal.blob_read import fetch_blob_bodies
+        from pypaimon.common.file_io import FileIO
+        from pypaimon.table.row.blob import BlobDescriptor
+
+        path = os.path.join(self.temp_dir, "array-blob.bin")
+        with open(path, "wb") as stream:
+            stream.write(b"0123456789")
+        file_io = FileIO.get("file://" + self.temp_dir, {})
+        descriptor = BlobDescriptor(path, 2, 3).serialize()
+        cells = {
+            "image": [descriptor],
+            "pages": [[descriptor, None, b"", b"inline", descriptor], [], 
None],
+            "assets": [[("cover", descriptor), ("missing", None)]],
+        }
+        with patch.object(file_io, "read_ranges_coalesced",
+                          wraps=file_io.read_ranges_coalesced) as read:
+            bodies = fetch_blob_bodies(
+                file_io, cells, list(cells), 2,
+                map_blob_cols=["assets"], array_blob_cols=["pages"])
+        self.assertEqual(1, read.call_count)
+        self.assertEqual({
+            "image": [b"234"],
+            "pages": [[b"234", None, b"", b"inline", b"234"], [], None],
+            "assets": [[("cover", b"234"), ("missing", None)]],
+        }, bodies)
+
     def test_fetch_bodies_rejects_unresolved_blob_view(self):
         from pypaimon.multimodal.query import ScanQuery
         from pypaimon.table.row.blob import BlobViewStruct

Reply via email to