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 38010b7e05 [python] Resolve ARRAY BLOB payloads in training windows 
(#10095)
38010b7e05 is described below

commit 38010b7e05f7c2f00e023e70e95d929e33fdbc43
Author: chaoyang <[email protected]>
AuthorDate: Thu Sep 24 14:38:21 2026 +0800

    [python] Resolve ARRAY BLOB payloads in training windows (#10095)
---
 .../pypaimon/multimodal/window_dataset.py          | 13 ++++--
 .../tests/contiguous_window_dataset_test.py        | 52 ++++++++++++++++++++++
 2 files changed, 62 insertions(+), 3 deletions(-)

diff --git a/paimon-python/pypaimon/multimodal/window_dataset.py 
b/paimon-python/pypaimon/multimodal/window_dataset.py
index ec5f639969..7b158eba83 100644
--- a/paimon-python/pypaimon/multimodal/window_dataset.py
+++ b/paimon-python/pypaimon/multimodal/window_dataset.py
@@ -36,7 +36,9 @@ from pypaimon.read.datasource.torch_dataset import (
     select_indexed_splits,
 )
 from pypaimon.read.query_auth_split import QueryAuthSplit
-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_file_type, is_map_blob_type,
+)
 from pypaimon.snapshot.time_travel_util import SCAN_KEYS
 from pypaimon.table.special_fields import SpecialFields
 from pypaimon.utils.range import Range
@@ -440,12 +442,16 @@ class _PinnedRowIdPlan:
         self._blob_columns = [
             field.name for field in table.fields
             if field.name in columns
-            and (is_blob_type(field.type) or is_map_blob_type(field.type))
+            and is_blob_file_type(field.type)
         ]
         self._map_blob_columns = {
             field.name for field in table.fields
             if field.name in self._blob_columns and 
is_map_blob_type(field.type)
         }
+        self._array_blob_columns = {
+            field.name for field in table.fields
+            if field.name in self._blob_columns and 
is_array_blob_type(field.type)
+        }
         blob_column_set = set(self._blob_columns)
         self._projection = (
             [name for name in columns if name not in blob_column_set]
@@ -497,7 +503,8 @@ class _PinnedRowIdPlan:
             arrow.select(self._blob_columns).to_pydict(),
             self._blob_columns,
             self._blob_parallelism,
-            self._map_blob_columns)
+            self._map_blob_columns,
+            self._array_blob_columns)
         blob_column_set = set(self._blob_columns)
         rows = arrow.select([
             name for name in arrow.column_names if name not in blob_column_set
diff --git a/paimon-python/pypaimon/tests/contiguous_window_dataset_test.py 
b/paimon-python/pypaimon/tests/contiguous_window_dataset_test.py
index c498cca3d0..0fd9d218cb 100644
--- a/paimon-python/pypaimon/tests/contiguous_window_dataset_test.py
+++ b/paimon-python/pypaimon/tests/contiguous_window_dataset_test.py
@@ -439,6 +439,58 @@ class ContiguousWindowDatasetTest(unittest.TestCase):
         self.assertEqual([b"scalar-0", b"scalar-1"], mixed["payload"])
         self.assertEqual(map_only["attachments"], mixed["attachments"])
 
+    def test_reads_array_blob_payloads_in_single_and_batched_windows(self):
+        pages = [[b"first", None, b"", b"first"], None, [], [b"last"]]
+        schema = pa.schema([
+            ("episode", pa.string()), ("step", pa.int32()),
+            ("pages", pa.list_(pa.large_binary())),
+            ("payload", pa.large_binary()),
+            ("attachments", pa.map_(pa.string(), pa.large_binary())),
+        ])
+        table = self.conn.create_table(
+            "array_blobs", schema=schema, options=_TABLE_OPTIONS)
+        table.add(pa.Table.from_pydict({
+            "episode": ["episode-a"] * 4, "step": [0, 1, 2, 3],
+            "pages": pages, "payload": [b"cover"] * 4,
+            "attachments": [[("thumb", b"thumbnail")]] * 4,
+        }, schema=schema))
+
+        for columns in (["pages"], ["pages", "payload", "attachments"]):
+            with self.subTest(columns=columns):
+                dataset = table.scan().to_contiguous_window_dataset(
+                    window_size=4, columns=columns,
+                    group_key="episode", order_key="step")
+                self.assertEqual(pages, dataset[0]["pages"])
+                self.assertEqual([pages, pages], [
+                    row["pages"] for row in dataset.__getitems__([0, 0])])
+                if "payload" in columns:
+                    self.assertEqual([b"cover"] * 4, dataset[0]["payload"])
+                    self.assertEqual(
+                        [[("thumb", b"thumbnail")]] * 4,
+                        dataset[0]["attachments"])
+
+        # Reordered/repeated offsets and endpoint padding retain cell 
structure.
+        dataset = table.scan().to_contiguous_window_dataset(
+            columns=["pages", "payload"], group_key="episode", 
order_key="step",
+            frame_offsets={"pages": [-1, 0, 2, 0]}, boundary="pad")
+        samples = dataset.__getitems__([3, 0, 3])
+        self.assertEqual([
+            [pages[2], pages[3], pages[3], pages[3]],
+            [pages[0], pages[0], pages[2], pages[0]],
+            [pages[2], pages[3], pages[3], pages[3]],
+        ], [sample["pages"] for sample in samples])
+        self.assertEqual([False, False, True, False],
+                         samples[0]["pages_is_pad"].tolist())
+        samples[0]["pages"][1].append(b"changed")
+        self.assertEqual([b"last"], samples[2]["pages"][1])
+
+        # Constant padding must not turn an empty ARRAY into a null cell.
+        padded = table.scan().to_contiguous_window_dataset(
+            columns=["pages"], group_key="episode", order_key="step",
+            frame_offsets={"pages": [-1, 0]}, boundary="pad",
+            pad_values={"pages": []})
+        self.assertEqual([[], pages[0]], padded[0]["pages"])
+
     def test_anchor_columns_read_only_the_window_anchor(self):
         table = self._table()
         with patch(

Reply via email to