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 9e306a4eef [python] Support encoding images as video (#10171)
9e306a4eef is described below

commit 9e306a4eefdd3a9bf910379a4f16bc0c81c39283
Author: Yann Byron <[email protected]>
AuthorDate: Fri Sep 25 16:19:09 2026 +0800

    [python] Support encoding images as video (#10171)
---
 .github/workflows/ci-python.yml                    |  2 +-
 docs/docs/pypaimon/video.md                        | 27 ++++++
 paimon-python/pypaimon/multimodal/table.py         | 48 +++++++++++
 paimon-python/pypaimon/multimodal/video.py         | 96 ++++++++++++++++++++-
 .../pypaimon/tests/multimodal_table_test.py        | 98 ++++++++++++++++++++++
 .../pypaimon/tests/multimodal_video_test.py        | 93 ++++++++++++++++++++
 paimon-python/setup.py                             |  6 ++
 7 files changed, 368 insertions(+), 2 deletions(-)

diff --git a/.github/workflows/ci-python.yml b/.github/workflows/ci-python.yml
index 892a155b5c..6c756f578e 100644
--- a/.github/workflows/ci-python.yml
+++ b/.github/workflows/ci-python.yml
@@ -147,7 +147,7 @@ jobs:
               python -c "import rosbags; print('rosbags installed')"
             fi
             if [[ "${{ matrix.python-version }}" == "3.10" ]]; then
-              python -m pip install './paimon-python[lerobot]'
+              python -m pip install './paimon-python[lerobot,video]'
               python -c "import datasets, lerobot; print('datasets', 
datasets.__version__, 'lerobot', lerobot.__version__)"
             fi
 
diff --git a/docs/docs/pypaimon/video.md b/docs/docs/pypaimon/video.md
index 881ce62fde..f01fe8c990 100644
--- a/docs/docs/pypaimon/video.md
+++ b/docs/docs/pypaimon/video.md
@@ -71,6 +71,33 @@ frames.add_videos([
 ])
 ```
 
+To encode an ordered sequence of image bytes or Pillow images before storing
+it, use Python 3.10 or newer, install the optional video dependencies, and call
+`add_images_as_video`:
+
+```shell
+pip install 'pypaimon[video]'
+```
+
+```python
+frames.add_images_as_video(
+    image_frames,
+    episode_44_rows,
+    fps=30,
+    codec="libx264",
+    pixel_format="yuv420p",
+    gop_size=2,
+    codec_options={"crf": "18", "preset": "fast"},
+)
+```
+
+The image and row counts must match, and every image must have the same
+dimensions. The container is MP4; codec-specific options are passed to PyAV.
+Encoding may be lossy and the original image payloads are not retained. A
+smaller GOP improves random frame access but usually increases file size.
+After encoding, Paimon copies the MP4 bytes into `.video` storage without
+transcoding them again.
+
 Each item may also be `(video, frame_rows, first_frame)`. Frame ordinals are
 generated consecutively from `first_frame`. If application semantics require
 PTS, wall-clock time, or a non-unit sampling map, retain that value in a normal
diff --git a/paimon-python/pypaimon/multimodal/table.py 
b/paimon-python/pypaimon/multimodal/table.py
index d988b91e07..e6f1bda9d3 100644
--- a/paimon-python/pypaimon/multimodal/table.py
+++ b/paimon-python/pypaimon/multimodal/table.py
@@ -133,6 +133,54 @@ class MultimodalTable:
             [(video, frames, first_frame)], video_column=video_column
         )
 
+    def add_images_as_video(
+            self,
+            images,
+            frames,
+            *,
+            fps,
+            codec,
+            pixel_format,
+            gop_size,
+            codec_options=None,
+            video_column=None):
+        """Encode ordered images as one MP4 and append its logical frame rows.
+
+        Images may be encoded image bytes or Pillow images. ``frames`` supplies
+        every table column except ``video_column`` and must contain one row per
+        image. Encoding finishes before the Paimon write starts.
+        """
+        import os
+        import tempfile
+
+        from pypaimon.multimodal.video import _encode_images_to_video
+
+        column = self._resolve_video_frame_column(video_column)
+        target_schema = _target_schema(self.raw_table)
+        non_video_schema = pa.schema([
+            field for field in target_schema if field.name != column
+        ])
+        frame_table = _to_arrow_table(frames, non_video_schema)
+        with tempfile.TemporaryDirectory(prefix="pypaimon_video_") as 
directory:
+            video_path = os.path.join(directory, "video.mp4")
+            image_count = _encode_images_to_video(
+                images,
+                video_path,
+                fps=fps,
+                codec=codec,
+                pixel_format=pixel_format,
+                gop_size=gop_size,
+                codec_options=codec_options,
+            )
+            if image_count != frame_table.num_rows:
+                raise ValueError(
+                    "Image count %d does not match frame row count %d."
+                    % (image_count, frame_table.num_rows)
+                )
+            return self.add_video(
+                video_path, frame_table, video_column=column
+            )
+
     def add_videos(self, videos, *, video_column=None):
         """Append several encoded videos with one writer and one commit.
 
diff --git a/paimon-python/pypaimon/multimodal/video.py 
b/paimon-python/pypaimon/multimodal/video.py
index ac2e7039e8..e105aeacf9 100644
--- a/paimon-python/pypaimon/multimodal/video.py
+++ b/paimon-python/pypaimon/multimodal/video.py
@@ -15,15 +15,109 @@
 # specific language governing permissions and limitations
 # under the License.
 
-"""PyTorch DataLoader helpers for descriptor-backed video frame rows."""
+"""Helpers for descriptor-backed video frame rows."""
 
+import io
 import os
 from collections import OrderedDict
 from collections.abc import Mapping
+from fractions import Fraction
+from itertools import chain
 
 from pypaimon.table.row.blob import Blob, VideoFrameDescriptor
 
 
+def _encode_images_to_video(
+        images,
+        output_path,
+        *,
+        fps,
+        codec,
+        pixel_format,
+        gop_size,
+        codec_options=None):
+    if isinstance(fps, bool) or not isinstance(fps, int) or fps <= 0:
+        raise ValueError("fps must be a positive int.")
+    if not isinstance(codec, str) or not codec:
+        raise ValueError("codec must be a non-empty string.")
+    if not isinstance(pixel_format, str) or not pixel_format:
+        raise ValueError("pixel_format must be a non-empty string.")
+    if isinstance(gop_size, bool) \
+            or not isinstance(gop_size, int) or gop_size <= 0:
+        raise ValueError("gop_size must be a positive int.")
+    if codec_options is not None and not isinstance(codec_options, Mapping):
+        raise ValueError("codec_options must be a mapping or None.")
+
+    try:
+        import av
+        from PIL import Image
+    except ImportError as error:
+        raise ImportError(
+            "Image-to-video encoding requires PyAV and Pillow; install "
+            "pypaimon[video]."
+        ) from error
+
+    iterator = iter(images)
+    try:
+        first = next(iterator)
+    except StopIteration as error:
+        raise ValueError("images must contain at least one frame.") from error
+
+    def open_image(value):
+        if isinstance(value, (bytes, bytearray, memoryview)):
+            image = Image.open(io.BytesIO(bytes(value)))
+        elif isinstance(value, Image.Image):
+            image = value
+        else:
+            raise TypeError(
+                "images must contain encoded image bytes or Pillow images."
+            )
+        image.load()
+        return image.convert("RGB")
+
+    first = open_image(first)
+    width, height = first.size
+    time_base = Fraction(1, fps)
+    count = 0
+    try:
+        with av.open(output_path, mode="w", format="mp4") as container:
+            stream = container.add_stream(
+                codec,
+                rate=fps,
+                options={
+                    str(key): str(value)
+                    for key, value in (codec_options or {}).items()
+                },
+            )
+            stream.width = width
+            stream.height = height
+            stream.pix_fmt = pixel_format
+            stream.gop_size = gop_size
+            stream.time_base = time_base
+            decoded = chain(
+                (first,), (open_image(value) for value in iterator)
+            )
+            for count, image in enumerate(decoded, start=1):
+                if image.size != (width, height):
+                    raise ValueError(
+                        "All images must have the same dimensions."
+                    )
+                frame = av.VideoFrame.from_image(image)
+                frame.pts = count - 1
+                frame.time_base = time_base
+                for packet in stream.encode(frame):
+                    container.mux(packet)
+            for packet in stream.encode():
+                container.mux(packet)
+    except BaseException:
+        try:
+            os.remove(output_path)
+        except FileNotFoundError:
+            pass
+        raise
+    return count
+
+
 class VideoFrameCollator:
     """Decode frame rows in a DataLoader worker while reusing video sessions.
 
diff --git a/paimon-python/pypaimon/tests/multimodal_table_test.py 
b/paimon-python/pypaimon/tests/multimodal_table_test.py
index 6551e743e0..be628983e7 100644
--- a/paimon-python/pypaimon/tests/multimodal_table_test.py
+++ b/paimon-python/pypaimon/tests/multimodal_table_test.py
@@ -161,6 +161,104 @@ class MultimodalTableTest(unittest.TestCase):
         _, bodies = table.scan().read_blobs("video", parallelism=2)
         self.assertEqual([video_bytes] * 3, bodies["video"])
 
+    def test_add_images_as_video_encodes_and_stores_frame_rows(self):
+        table = self.conn.create_table(
+            "image_video_frames",
+            schema=_schema({
+                "episode_id": pa.int64(),
+                "video": pa.large_binary(),
+            }),
+            options=dict(_PARQUET_OPTIONS, **{
+                "video-frame-field": "video",
+                "blob-as-descriptor": "true",
+            }),
+        )
+        calls = []
+
+        def encode(images, output_path, **options):
+            images = list(images)
+            calls.append((images, options))
+            with open(output_path, "wb") as output:
+                output.write(b"encoded-video")
+            return len(images)
+
+        with patch(
+                "pypaimon.multimodal.video._encode_images_to_video",
+                encode):
+            table.add_images_as_video(
+                [b"frame-0", b"frame-1"],
+                [{"episode_id": 42}, {"episode_id": 42}],
+                fps=30,
+                codec="mpeg4",
+                pixel_format="yuv420p",
+                gop_size=2,
+                codec_options={"qscale": "3"},
+            )
+
+        self.assertEqual(
+            [
+                (
+                    [b"frame-0", b"frame-1"],
+                    {
+                        "fps": 30,
+                        "codec": "mpeg4",
+                        "pixel_format": "yuv420p",
+                        "gop_size": 2,
+                        "codec_options": {"qscale": "3"},
+                    },
+                )
+            ],
+            calls,
+        )
+        rows = table.scan().select(["episode_id", "video"]).to_list()
+        descriptors = [
+            pmm.VideoFrameDescriptor.deserialize(row["video"])
+            for row in rows
+        ]
+        self.assertEqual([0, 1], [value.frame_index for value in descriptors])
+        _, bodies = table.scan().read_blobs("video")
+        self.assertEqual([b"encoded-video"] * 2, bodies["video"])
+
+    def test_add_images_as_video_rejects_frame_count_mismatch(self):
+        table = self.conn.create_table(
+            "mismatched_image_video_frames",
+            schema=_schema({
+                "episode_id": pa.int64(),
+                "video": pa.large_binary(),
+            }),
+            options=dict(_PARQUET_OPTIONS, **{
+                "video-frame-field": "video",
+                "blob-as-descriptor": "true",
+            }),
+        )
+        output_paths = []
+
+        def encode(unused_images, output_path, **unused_options):
+            output_paths.append(output_path)
+            with open(output_path, "wb") as output:
+                output.write(b"one-frame-video")
+            return 1
+
+        with patch(
+                "pypaimon.multimodal.video._encode_images_to_video",
+                encode):
+            with self.assertRaisesRegex(
+                    ValueError,
+                    "Image count 1 does not match frame row count 2"):
+                table.add_images_as_video(
+                    [b"frame-0"],
+                    [{"episode_id": 42}, {"episode_id": 42}],
+                    fps=30,
+                    codec="mpeg4",
+                    pixel_format="yuv420p",
+                    gop_size=2,
+                )
+
+        self.assertFalse(os.path.exists(output_paths[0]))
+        self.assertIsNone(
+            table.raw_table.snapshot_manager().get_latest_snapshot()
+        )
+
     def test_add_videos_packs_multiple_videos_in_one_commit(self):
         from pypaimon.table.row.blob import Blob, VideoFrameDescriptor
 
diff --git a/paimon-python/pypaimon/tests/multimodal_video_test.py 
b/paimon-python/pypaimon/tests/multimodal_video_test.py
index 918e917a32..750f6f142f 100644
--- a/paimon-python/pypaimon/tests/multimodal_video_test.py
+++ b/paimon-python/pypaimon/tests/multimodal_video_test.py
@@ -50,6 +50,99 @@ class _FailingCloseDecoder(_Decoder):
         raise RuntimeError("decoder close failed")
 
 
+class ImageVideoEncoderTest(unittest.TestCase):
+
+    def setUp(self):
+        self.temp_dir = tempfile.TemporaryDirectory()
+
+    def tearDown(self):
+        self.temp_dir.cleanup()
+
+    def test_encodes_ordered_image_bytes_as_video(self):
+        try:
+            import av
+            from PIL import Image
+        except ImportError:
+            self.skipTest("PyAV and Pillow are required for video encoding")
+        from pypaimon.multimodal.video import _encode_images_to_video
+
+        images = []
+        for color in ((255, 0, 0), (0, 255, 0)):
+            output = io.BytesIO()
+            Image.new("RGB", (8, 6), color).save(output, format="PNG")
+            images.append(output.getvalue())
+        video_path = os.path.join(self.temp_dir.name, "encoded.mp4")
+
+        count = _encode_images_to_video(
+            images,
+            video_path,
+            fps=10,
+            codec="mpeg4",
+            pixel_format="yuv420p",
+            gop_size=1,
+            codec_options={"qscale": "3"},
+        )
+
+        with av.open(video_path) as container:
+            stream = container.streams.video[0]
+            frames = list(container.decode(stream))
+            size = stream.width, stream.height
+        self.assertEqual(2, count)
+        self.assertEqual(2, len(frames))
+        self.assertEqual((8, 6), size)
+
+    def test_image_video_encoder_rejects_invalid_required_options(self):
+        from pypaimon.multimodal.video import _encode_images_to_video
+
+        defaults = {
+            "fps": 10,
+            "codec": "mpeg4",
+            "pixel_format": "yuv420p",
+            "gop_size": 1,
+        }
+        cases = [
+            ({"fps": True}, "fps must be a positive int"),
+            ({"codec": ""}, "codec must be a non-empty string"),
+            ({"pixel_format": ""},
+             "pixel_format must be a non-empty string"),
+            ({"gop_size": 0}, "gop_size must be a positive int"),
+            ({"codec_options": []}, "codec_options must be a mapping"),
+        ]
+        for overrides, message in cases:
+            options = dict(defaults)
+            options.update(overrides)
+            with self.subTest(options=options):
+                with self.assertRaisesRegex(ValueError, message):
+                    _encode_images_to_video(
+                        [],
+                        os.path.join(self.temp_dir.name, "invalid.mp4"),
+                        **options,
+                    )
+
+    def 
test_image_video_encoder_rejects_mixed_dimensions_and_cleans_output(self):
+        try:
+            from PIL import Image
+            import av  # noqa: F401
+        except ImportError:
+            self.skipTest("PyAV and Pillow are required for video encoding")
+        from pypaimon.multimodal.video import _encode_images_to_video
+
+        video_path = os.path.join(self.temp_dir.name, "mixed.mp4")
+        with self.assertRaisesRegex(ValueError, "same dimensions"):
+            _encode_images_to_video(
+                [
+                    Image.new("RGB", (8, 6), "red"),
+                    Image.new("RGB", (10, 6), "green"),
+                ],
+                video_path,
+                fps=10,
+                codec="mpeg4",
+                pixel_format="yuv420p",
+                gop_size=1,
+            )
+        self.assertFalse(os.path.exists(video_path))
+
+
 class VideoFrameCollatorTest(unittest.TestCase):
 
     def setUp(self):
diff --git a/paimon-python/setup.py b/paimon-python/setup.py
index 7d732167b8..47e600343e 100644
--- a/paimon-python/setup.py
+++ b/paimon-python/setup.py
@@ -217,6 +217,11 @@ def read_requirements():
 
 install_requires = read_requirements()
 
+VIDEO_DEPENDENCIES = [
+    'av>=12,<19; python_version>="3.10"',
+    'Pillow; python_version>="3.10"',
+]
+
 LEROBOT_DEPENDENCIES = [
     # datasets 4.1+ may select PyArrow 21+, while PyPaimon currently
     # supports PyArrow <20. Pandas 2.2.2+ supports NumPy 2.x selected
@@ -249,6 +254,7 @@ setup(
         ],
     },
     extras_require={
+        'video': VIDEO_DEPENDENCIES,
         'hdf5': [
             # HDF5 loading is explicitly guarded and documented as Python 3.8+.
             'h5py>=3,<4; python_version>="3.8"',

Reply via email to