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(