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 5889157343 [python] Support ARRAY and MAP BLOBs in Ray batch 
processing (#10060)
5889157343 is described below

commit 5889157343341ff136fd8aa64c099fe02e76699f
Author: chaoyang <[email protected]>
AuthorDate: Tue Sep 22 10:43:59 2026 +0800

    [python] Support ARRAY and MAP BLOBs in Ray batch processing (#10060)
---
 docs/docs/pypaimon/multimodal-reading.md           |   8 +
 paimon-python/pypaimon/multimodal/query.py         |  11 +-
 paimon-python/pypaimon/multimodal/table.py         |  10 +-
 paimon-python/pypaimon/ray/ray_paimon.py           |  84 ++++++++--
 .../pypaimon/tests/multimodal_table_test.py        | 175 +++++++++++++++++++++
 5 files changed, 270 insertions(+), 18 deletions(-)

diff --git a/docs/docs/pypaimon/multimodal-reading.md 
b/docs/docs/pypaimon/multimodal-reading.md
index eb0a60239e..fd3685ba27 100644
--- a/docs/docs/pypaimon/multimodal-reading.md
+++ b/docs/docs/pypaimon/multimodal-reading.md
@@ -324,6 +324,14 @@ values. Read those columns with `to_torch()` instead.
 For larger jobs, read descriptors with `to_ray()`, then fetch and process BLOB
 bytes on Ray workers with `map_with_blobs`.
 
+Scalar, ARRAY, and MAP BLOB columns are supported. In the callback's `blobs`
+dictionary, each column contains row-aligned values: bytes for scalar BLOBs,
+lists for ARRAY BLOBs, and lists of key-value pairs for MAP BLOBs. Null cells,
+null elements, empty containers, and element order are preserved. All source
+BLOB columns are excluded from `scalar_batch`, including unrequested ones.
+After Ray transformations, use `table.map_with_blobs(...)` to supply the source
+table's BLOB column information even if a transform changes Arrow nested types.
+
 ```python
 import ray
 import pyarrow as pa
diff --git a/paimon-python/pypaimon/multimodal/query.py 
b/paimon-python/pypaimon/multimodal/query.py
index fdcd33c3a0..17ea818eef 100644
--- a/paimon-python/pypaimon/multimodal/query.py
+++ b/paimon-python/pypaimon/multimodal/query.py
@@ -275,7 +275,10 @@ class ScanQuery:
                 fn_kwargs={"columns": visible_columns},
                 batch_format="pyarrow")
         setattr(ds, "_paimon_blob_file_io", file_io)
-        setattr(ds, "_paimon_blob_columns", self._all_blob_columns())
+        setattr(ds, "_paimon_blob_columns", self._readable_blob_columns())
+        map_blob_cols, array_blob_cols = self._nested_blob_columns()
+        setattr(ds, "_paimon_map_blob_columns", map_blob_cols)
+        setattr(ds, "_paimon_array_blob_columns", array_blob_cols)
         return ds
 
     def read_blobs(
@@ -396,12 +399,6 @@ class ScanQuery:
             [field.name for field in self._table.fields if 
is_array_blob_type(field.type)],
         )
 
-    def _all_blob_columns(self) -> List[str]:
-        return [
-            field.name for field in self._table.fields
-            if is_blob_type(field.type)
-        ]
-
     def _readable_blob_columns(self) -> List[str]:
         return [
             field.name for field in self._table.fields
diff --git a/paimon-python/pypaimon/multimodal/table.py 
b/paimon-python/pypaimon/multimodal/table.py
index d7b8706b51..4931da1dd1 100644
--- a/paimon-python/pypaimon/multimodal/table.py
+++ b/paimon-python/pypaimon/multimodal/table.py
@@ -30,7 +30,9 @@ from pypaimon.multimodal.query import (
     TextQuery,
     VectorQuery,
 )
-from pypaimon.schema.data_types import PyarrowFieldParser, is_blob_type
+from pypaimon.schema.data_types import (
+    PyarrowFieldParser, is_array_blob_type, is_blob_file_type, 
is_map_blob_type,
+)
 from pypaimon.table.data_evolution_merge_into import (
     WhenMatched,
     WhenNotMatched,
@@ -355,6 +357,10 @@ class MultimodalTable:
             fn,
             file_io=self.raw_table.file_io,
             all_blob_columns=_blob_columns(self.raw_table),
+            map_blob_columns=[field.name for field in self.raw_table.fields
+                              if is_map_blob_type(field.type)],
+            array_blob_columns=[field.name for field in self.raw_table.fields
+                                if is_array_blob_type(field.type)],
             **kwargs,
         )
 
@@ -569,7 +575,7 @@ class _MergeBuilder:
 def _blob_columns(table):
     return tuple(
         field.name for field in table.fields
-        if is_blob_type(field.type)
+        if is_blob_file_type(field.type)
     )
 
 
diff --git a/paimon-python/pypaimon/ray/ray_paimon.py 
b/paimon-python/pypaimon/ray/ray_paimon.py
index 460344fb6e..58358fe81c 100644
--- a/paimon-python/pypaimon/ray/ray_paimon.py
+++ b/paimon-python/pypaimon/ray/ray_paimon.py
@@ -147,6 +147,8 @@ def map_with_blobs(
     *,
     file_io=None,
     all_blob_columns=None,
+    map_blob_columns=None,
+    array_blob_columns=None,
     parallelism: int = 64,
     batch_size: Optional[int] = 1024,
     fn_kwargs: Optional[Dict[str, Any]] = None,
@@ -156,10 +158,15 @@ def map_with_blobs(
     """Fetch BLOB payloads in Ray batches and call ``fn``.
 
     ``fn(scalar_batch, blobs, **fn_kwargs)`` receives a ``pyarrow.Table`` of
-    non-BLOB columns and a row-aligned ``dict`` of BLOB bytes. Return a small
-    Ray-compatible batch; for side-effect-only work, return an empty
+    non-BLOB columns and a row-aligned ``dict`` of BLOB bytes. ARRAY BLOB cells
+    are lists and MAP BLOB cells are lists of key-value pairs; null cells and
+    elements are preserved. Return a small Ray-compatible batch; for
+    side-effect-only work, return an empty
     ``pyarrow.Table`` instead of ``None``. Call this directly on
     ``scan().to_ray()`` output, or pass ``file_io`` and ``all_blob_columns``.
+    Supply ``map_blob_columns`` and ``array_blob_columns`` when transforms
+    erase the source metadata and Arrow nested types. The table method
+    supplies this information automatically.
     Tune ``batch_size`` for BLOB size and worker memory.
     """
     _require_ray_data()
@@ -209,12 +216,23 @@ def map_with_blobs(
     if invalid:
         raise ValueError("Column {!r} is not a BLOB 
column.".format(invalid[0]))
 
+    if map_blob_columns is None:
+        map_blob_columns = getattr(dataset, "_paimon_map_blob_columns", ())
+    if array_blob_columns is None:
+        array_blob_columns = getattr(dataset, "_paimon_array_blob_columns", ())
+    map_blob_columns, array_blob_columns = set(map_blob_columns), 
set(array_blob_columns)
+    if (map_blob_columns & array_blob_columns
+            or (map_blob_columns | array_blob_columns) - all_blob):
+        raise ValueError("Nested BLOB columns must be disjoint subsets of 
all_blob_columns.")
+
     return dataset.map_batches(
         _map_blob_batch,
         fn_kwargs={
             "file_io": resolved_file_io,
             "blob_cols": blob_cols,
             "all_blob_cols": list(all_blob_cols),
+            "map_blob_cols": list(map_blob_columns),
+            "array_blob_cols": list(array_blob_columns),
             "parallelism": parallelism,
             "fn": fn,
             "fn_kwargs": dict(fn_kwargs or {}),
@@ -233,8 +251,11 @@ def _set_map_batches_remote_args(dataset, kwargs, 
ray_remote_args):
 
 
 def _map_blob_batch(
-        batch, file_io, blob_cols, all_blob_cols, parallelism, fn, fn_kwargs):
-    from pypaimon.multimodal.blob_read import fetch_blob_bodies
+        batch, file_io, blob_cols, all_blob_cols, parallelism, fn, fn_kwargs,
+        map_blob_cols=(), array_blob_cols=()):
+    import pyarrow as pa
+
+    from pypaimon.multimodal.blob_read import _map_entries, fetch_blob_bodies
 
     missing = [name for name in blob_cols if name not in batch.schema.names]
     if missing:
@@ -251,8 +272,28 @@ def _map_blob_batch(
             "table.map_with_blobs() in a separate pass, or drop it before "
             "mapping.".format(unknown[0]))
 
+    map_blob_cols = set(map_blob_cols) | {
+        name for name in blob_cols if 
pa.types.is_map(batch.schema.field(name).type)}
+    array_blob_cols = set(array_blob_cols) | {
+        name for name in blob_cols
+        if (pa.types.is_list(batch.schema.field(name).type)
+            or pa.types.is_large_list(batch.schema.field(name).type)
+            or pa.types.is_fixed_size_list(batch.schema.field(name).type))
+    }
+    data = batch.select(blob_cols).to_pydict()
+    # Ray row transforms may infer MAP pairs as list<list<string>> when the
+    # payload bytes are valid UTF-8. Restore the values without changing keys.
+    for name in map_blob_cols.intersection(blob_cols):
+        data[name] = [
+            None if row is None else [
+                (key, value.encode("utf-8") if isinstance(value, str) else 
value)
+                for key, value in _map_entries(row)
+            ]
+            for row in data[name]
+        ]
     bodies = fetch_blob_bodies(
-        file_io, batch.select(blob_cols).to_pydict(), blob_cols, parallelism)
+        file_io, data, blob_cols, parallelism,
+        map_blob_cols, array_blob_cols)
     result = fn(batch.select(scalar_cols), bodies, **fn_kwargs)
     if result is None:
         raise ValueError(
@@ -268,17 +309,42 @@ def _unknown_blob_descriptor_columns(batch, scalar_cols):
         if _looks_like_blob_descriptor(batch.column(name))]
 
 
-def _looks_like_blob_descriptor(column):
+def _looks_like_blob_descriptor(column, nested=False):
     import pyarrow as pa
     from pypaimon.table.row.blob import BlobDescriptorSerde
 
+    chunks = getattr(column, "chunks", [column])
+    if (pa.types.is_list(column.type) or pa.types.is_large_list(column.type)
+            or pa.types.is_fixed_size_list(column.type)):
+        return any(_looks_like_blob_descriptor(chunk.flatten(), nested=True) 
for chunk in chunks)
+    if pa.types.is_map(column.type):
+        return any(
+            _looks_like_blob_descriptor(value.values.field(1), nested=True)
+            for value in column if value.is_valid)
+    if isinstance(column.type, pa.ExtensionType):
+        return any(_contains_blob_descriptor(value.as_py()) for value in 
column if value.is_valid)
+    if nested and (pa.types.is_string(column.type) or 
pa.types.is_large_string(column.type)):
+        return any(_contains_blob_descriptor(value.as_py()) for value in 
column if value.is_valid)
     if not (pa.types.is_binary(column.type) or 
pa.types.is_large_binary(column.type)):
         return False
-    chunks = getattr(column, "chunks", None) or [column]
     for chunk in chunks:
         for value in chunk:
-            if value.is_valid:
-                return BlobDescriptorSerde.is_descriptor(value.as_py())
+            if value.is_valid and 
BlobDescriptorSerde.is_descriptor(value.as_py()):
+                return True
+    return False
+
+
+def _contains_blob_descriptor(value):
+    from pypaimon.table.row.blob import BlobDescriptorSerde
+
+    if isinstance(value, str):
+        value = value.encode("utf-8")
+    if isinstance(value, (bytes, bytearray, memoryview)):
+        return BlobDescriptorSerde.is_descriptor(bytes(value))
+    if isinstance(value, dict):
+        return any(_contains_blob_descriptor(item) for item in value.values())
+    if isinstance(value, (list, tuple)):
+        return any(_contains_blob_descriptor(item) for item in value)
     return False
 
 
diff --git a/paimon-python/pypaimon/tests/multimodal_table_test.py 
b/paimon-python/pypaimon/tests/multimodal_table_test.py
index 4d7b2c2b33..b5e04a7238 100644
--- a/paimon-python/pypaimon/tests/multimodal_table_test.py
+++ b/paimon-python/pypaimon/tests/multimodal_table_test.py
@@ -1750,6 +1750,181 @@ class MultimodalTableTest(unittest.TestCase):
             if started_ray:
                 ray.shutdown()
 
+    @unittest.skipIf(ray is None, "ray is not installed")
+    def test_scan_to_ray_map_with_nested_blobs(self):
+        from pypaimon.ray import map_with_blobs
+
+        schema = _schema({
+            "id": pa.int32(), "image": pa.large_binary(),
+            "pages": pa.list_(pa.large_binary()),
+            "assets": pa.map_(pa.string(), pa.large_binary()),
+        })
+        table = self.conn.create_table(
+            "ray_nested_blobs", schema=schema, options=_PARQUET_OPTIONS)
+        rows = [
+            {"id": 1, "image": b"preview", "pages": [b"a", None, b"", b"b"],
+             "assets": [("thumb", b"thumb"), ("missing", None), ("empty", 
b"")]},
+            {"id": 2, "image": None, "pages": None, "assets": None},
+            {"id": 3, "image": b"preview", "pages": [], "assets": []},
+            {"id": 4, "image": b"preview", "pages": [b"d"], "assets": 
[("thumb", b"last")]},
+        ]
+        table.add(pa.Table.from_pylist(rows, schema=schema))
+        output_schema = pa.schema([schema.field(name) for name in ("id", 
"pages", "assets")])
+
+        def collect(scalar, blobs):
+            assert scalar.column_names == ["id"]
+            return pa.Table.from_pydict(dict(id=scalar["id"], **blobs), 
schema=output_schema)
+
+        started_ray = not ray.is_initialized()
+        if started_ray:
+            ray.init(ignore_reinit_error=True, num_cpus=2)
+        try:
+            dataset = table.scan().to_ray(concurrency=1, override_num_blocks=1)
+            for maps, arrays in ((["assets"], ["assets"]), (["unknown"], [])):
+                with self.assertRaisesRegex(ValueError, "disjoint subsets"):
+                    map_with_blobs(dataset, ["pages", "assets"], collect,
+                                   map_blob_columns=maps, 
array_blob_columns=arrays)
+            expected = [{name: row[name] for name in output_schema.names} for 
row in rows]
+            for use_table in (False, True):
+                with self.subTest(use_table=use_table):
+                    if use_table:
+                        # Ray transformations do not preserve the source 
metadata.
+                        result = table.map_with_blobs(
+                            dataset.filter(lambda row: row["id"] > 0),
+                            ["pages", "assets"], collect, batch_size=2)
+                    else:
+                        result = map_with_blobs(
+                            dataset, ["pages", "assets"], collect, 
batch_size=2)
+                    actual = [row for batch in 
result.iter_batches(batch_format="pyarrow")
+                              for row in batch.to_pylist()]
+                    self.assertEqual(expected, sorted(actual, key=lambda row: 
row["id"]))
+
+            projected = table.scan().select(["id", "pages"]).to_ray()
+            result = table.map_with_blobs(
+                projected, "pages", lambda scalar, blobs: scalar, batch_size=2)
+            self.assertEqual([1, 2, 3, 4], sorted(row["id"] for row in 
result.take_all()))
+            empty = table.scan().where("id < 0").to_ray()
+            self.assertEqual([], map_with_blobs(empty, ["pages", "assets"], 
collect).take_all())
+        finally:
+            if started_ray:
+                ray.shutdown()
+
+    def test_map_with_blobs_rejects_foreign_nested_descriptors(self):
+        from pypaimon.filesystem.local_file_io import LocalFileIO
+        from pypaimon.ray.ray_paimon import _map_blob_batch
+        from pypaimon.table.row.blob import BlobDescriptor
+
+        descriptor = BlobDescriptor("oss://other-table/blob", 0, 1).serialize()
+        for arrow_type, values in (
+                (pa.list_(pa.large_binary()), [None, [], [b"inline", 
descriptor]]),
+                (pa.map_(pa.string(), pa.large_binary()),
+                 [None, [], [("inline", b"inline"), ("foreign", descriptor)]]),
+                (pa.list_(pa.list_(pa.string())),
+                 [None, [], [("inline", "inline"), ("foreign", 
descriptor.decode("utf-8"))]])):
+            with self.subTest(arrow_type=arrow_type):
+                batch = pa.table({
+                    "image": [b"inline"] * 3,
+                    "text": [descriptor.decode("utf-8")] * 3,
+                    "foreign": pa.array(values, type=arrow_type),
+                })
+                with self.assertRaisesRegex(ValueError, "does not own"):
+                    _map_blob_batch(batch, None, ["image"], ["image"], 1, 
None, {})
+                # A descriptor outside the current slice is not a foreign 
column value.
+                _map_blob_batch(batch.slice(0, 2), LocalFileIO(), ["image"], 
["image"], 1,
+                                lambda scalar, blobs: scalar, {})
+                empty = pa.table({
+                    "image": pa.chunked_array([], type=pa.large_binary()),
+                    "foreign": pa.chunked_array([], type=arrow_type),
+                })
+                result = _map_blob_batch(empty, LocalFileIO(), ["image"], 
["image"], 1,
+                                         lambda scalar, blobs: scalar, {})
+                self.assertEqual(0, result.num_rows)
+
+    def test_map_with_blobs_restores_utf8_map_values(self):
+        from unittest.mock import Mock
+
+        from pypaimon.ray.ray_paimon import _map_blob_batch
+        from pypaimon.table.row.blob import BlobDescriptor
+
+        descriptor = BlobDescriptor("test://blob", 0, 4).serialize()
+        entries = [("blob", descriptor.decode("utf-8")), ("text", 
"\u4f60\u597d"),
+                   ("empty", ""), ("missing", None), ("blob", 
descriptor.decode("utf-8"))]
+        batch = pa.table({"assets": pa.array([None, [], entries],
+                                             
type=pa.list_(pa.list_(pa.string())))})
+        file_io = Mock()
+        file_io.read_ranges_coalesced.return_value = [b"body", None, None, 
None, b"body"]
+        actual = _map_blob_batch(
+            batch, file_io, ["assets"], ["assets"], 2, lambda scalar, blobs: 
blobs, {},
+            map_blob_cols=["assets"])
+        self.assertEqual({"assets": [None, [], [
+            ("blob", b"body"), ("text", "\u4f60\u597d".encode("utf-8")),
+            ("empty", b""), ("missing", None), ("blob", b"body"),
+        ]]}, actual)
+        file_io.read_ranges_coalesced.assert_called_once_with(
+            [("test://blob", 0, 4), None, None, None, ("test://blob", 0, 4)], 
2)
+
+    @unittest.skipIf(ray is None, "ray is not installed")
+    def test_ray_filter_preserves_map_blob_bytes(self):
+        from pypaimon.filesystem.local_file_io import LocalFileIO
+        from pypaimon.ray import map_with_blobs
+
+        schema = _schema({"id": pa.int32(), "assets": pa.map_(pa.string(), 
pa.large_binary())})
+
+        def collect(scalar, blobs):
+            return pa.Table.from_pydict(dict(id=scalar["id"], **blobs), 
schema=schema)
+
+        started_ray = not ray.is_initialized()
+        if started_ray:
+            ray.init(ignore_reinit_error=True, num_cpus=2)
+        try:
+            # Exercise string inference and object fallback independently of 
descriptor URI length.
+            for payload in ("\u4f60\u597d\0".encode("utf-8"), b"\xff\0"):
+                with self.subTest(payload=payload):
+                    rows = [
+                        {"id": 1, "assets": None}, {"id": 2, "assets": []},
+                        {"id": 3, "assets": [("body", payload), ("empty", 
b""), ("missing", None)]},
+                    ]
+                    dataset = ray.data.from_arrow(pa.Table.from_pylist(rows, 
schema=schema))
+                    result = map_with_blobs(
+                        dataset.filter(lambda row: row["id"] > 0), "assets", 
collect,
+                        file_io=LocalFileIO(), all_blob_columns=["assets"], 
map_blob_columns=["assets"],
+                        batch_size=2)
+                    actual = [row for batch in 
result.iter_batches(batch_format="pyarrow")
+                              for row in batch.to_pylist()]
+                    self.assertEqual(rows, sorted(actual, key=lambda row: 
row["id"]))
+        finally:
+            if started_ray:
+                ray.shutdown()
+
+    @unittest.skipIf(ray is None, "ray is not installed")
+    def test_map_with_blobs_python_object_columns(self):
+        from pypaimon.filesystem.local_file_io import LocalFileIO
+        from pypaimon.ray.ray_paimon import _map_blob_batch
+        from pypaimon.table.row.blob import BlobDescriptor
+
+        try:
+            from ray.data.extensions import ArrowPythonObjectArray
+        except ImportError:
+            self.skipTest("Ray has no Python object extension type")
+
+        expected = {
+            "assets": [None, [], [("image", b"body"), ("missing", None)]],
+            "pages": [None, [], [b"page", None]],
+        }
+        batch = pa.table({name: ArrowPythonObjectArray.from_objects(values)
+                          for name, values in expected.items()})
+        actual = _map_blob_batch(
+            batch, LocalFileIO(), ["assets", "pages"], ["assets", "pages"], 1,
+            lambda scalar, blobs: blobs, {}, map_blob_cols=["assets"], 
array_blob_cols=["pages"])
+        self.assertEqual(expected, actual)
+
+        foreign = ArrowPythonObjectArray.from_objects([
+            None, [], [("inline", b"body"), ("foreign", 
BlobDescriptor("oss://other/blob", 0, 1).serialize())],
+        ])
+        with self.assertRaisesRegex(ValueError, "does not own"):
+            _map_blob_batch(batch.append_column("foreign", foreign), None,
+                            ["assets"], ["assets", "pages"], 1, None, {}, 
map_blob_cols=["assets"])
+
     @unittest.skipIf(ray is None, "ray is not installed")
     def test_scan_to_ray_map_with_blobs_guards(self):
         started_ray = False

Reply via email to