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