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