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"',