Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 22 additions & 1 deletion docs/docs/concepts/spec/fileformat.md
Original file line number Diff line number Diff line change
Expand Up @@ -413,15 +413,22 @@ the ordinary BLOB entry header, length trailer, or per-entry CRC:
+----------------------------+
| ... |
+----------------------------+
| Keyframe Index 1 | Video metadata ranges and compressed keyframe entries
+----------------------------+
| Keyframe Index 2 |
+----------------------------+
| Physical Length Index | Delta-Varint video lengths
+----------------------------+
| Keyframe-Index Length Index | Delta-Varint keyframe-index lengths
+----------------------------+
| Run Length Index | Delta-Varint logical row counts
+----------------------------+
| Run Reference Index | Delta-Varint physical video ordinals
+----------------------------+
| Run First-Frame Index | Delta-Varint frame ordinals
+----------------------------+
| Physical Index Length | 4 bytes (Little Endian)
| Keyframe Length-Index Size | 4 bytes (Little Endian)
| Run-Length Index Length | 4 bytes (Little Endian)
| Run-Reference Index Length | 4 bytes (Little Endian)
| First-Frame Index Length | 4 bytes (Little Endian)
Expand All @@ -434,7 +441,19 @@ The run arrays have equal element counts. A non-negative run reference is an ord
physical length index. For logical row `r` in a run beginning at logical row `s`, the returned
`VideoFrameDescriptor` identifies the referenced raw video range and frame ordinal
`run_first_frame + (r - s)`. `-1` is a NULL run and `-2` is a data-evolution placeholder run.
Non-negative runs have fixed frame stride one in version 1; a discontinuity starts another run.
Non-negative runs have fixed frame stride one; a discontinuity starts another run. Each physical
video may store one sparse keyframe index; a zero length selects the scan fallback.

A keyframe-index block has a 17-byte header: version (`1`, uint8), magic (`0x564944454F4B4649`, uint64),
metadata-range count (uint32), and keyframe count (uint32). It then stores metadata `(offset,
length)` pairs (two int64 values) and zlib-compressed `(frame ordinal, PTS, packet position)`
keyframe entries (three int64 values). Offsets are relative to the encoded video; writers reject
out-of-range values. All numeric fields are little endian. One block is limited to 65,536
metadata ranges, 65,536 keyframes, and 16 MiB; all blocks in one file are limited to 64 MiB.

The index covers the first video stream. Its time base remains in the video. A reader fetches the
metadata and target GOP, seeks by PTS, and decodes forward in presentation order. It may include the
following GOP for reordered frames.

The serialized `VideoFrameDescriptor` stored in an Arrow/data-file cell has its own versioned
wire layout. All numeric values are little endian:
Expand All @@ -448,6 +467,8 @@ wire layout. All numeric values are little endian:
| Offset | 8 bytes | Start of the complete encoded-video payload |
| Length | 8 bytes | Encoded-video payload length |
| Frame index | 8 bytes | Zero-based presentation-order frame ordinal |
| Keyframe-index offset | 8 bytes | Index offset in the `.video` file, or `-1` |
| Keyframe-index length | 8 bytes | Index length, or `0` |

Descriptor bytes are independently versioned from the `.video` container. Java and Python share
canonical descriptor and container fixtures to keep both implementations byte-compatible.
Expand Down
37 changes: 25 additions & 12 deletions docs/docs/multimodal-table/video.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -81,12 +81,28 @@ to write complete videos and their logical frame rows.

## Descriptors and Physical Layout

Each logical value is a `VideoFrameDescriptor`: its URI range identifies one complete encoded
video and its frame ordinal selects a frame inside that video. On write, a `.video` data region
concatenates raw video payloads without ordinary BLOB entry wrappers. Four embedded delta-varint
indexes record physical video lengths, logical run lengths, run-to-video references, and the first
frame ordinal of each run. Consecutive frames therefore need one run entry rather than one index
entry per row. NULL and data-evolution placeholders use negative run references.
Each `VideoFrameDescriptor` identifies an encoded video range and a frame ordinal. A `.video` file
packs raw videos, optional keyframe indexes, and delta-varint row mappings. Consecutive frames share
one run; negative references represent NULL or data-evolution placeholders.

<pre>
.video (version 1, in file order)
├─ Complete encoded videos A, B, ...
├─ Keyframe index blocks A, B, ...
│ ├─ Video metadata ranges: offset and length within the encoded video
│ └─ Compressed keyframe entries: frame ordinal, PTS, and packet byte position
├─ Metadata indexes
│ ├─ Video payload lengths
│ ├─ Keyframe index block lengths
│ └─ Run lengths, video references, and starting frame indices
└─ Footer: five 4-byte index sizes, 4-byte magic, and 1-byte version
</pre>

Each video index stores initialization ranges such as MP4 `moov` and keyframe positions. A cold
reader prefetches these ranges and a bounded GOP window, then fetches uncached decoder reads on
demand. Zero length selects the scan fallback. PyPaimon generates indexes for supported ISO BMFF
videos and preserves them during rewrites. See the
[file format specification](../concepts/spec/fileformat#video).

## Reuse, Rolling, and Compaction

Expand All @@ -101,10 +117,9 @@ BLOB, and vector rolling is deferred until the next payload boundary and all act
together. The resulting file group remains row-aligned and a single large episode may exceed the
configured target.

When BLOB compaction is enabled, it byte-copies the complete encoded-video ranges into a new self-contained `.video` pack
and rebuilds the embedded indexes; it does not decode or re-encode frames. The aligned normal data
file contains only application columns such as `episode_id`, state, and action. Paimon stores the
frame mapping in the `.video` descriptor/index path.
When BLOB compaction is enabled, it copies videos and keyframe indexes into a new self-contained
`.video` pack without decoding frames. The aligned normal data file contains only application
columns such as `episode_id`, state, and action.
`blob-compaction.enabled` defaults to `false`; ordinary normal-file compaction can
leave the video packs unchanged. See [Data Evolution Maintenance](./data-evolution-maintenance#ordinary-compaction).

Expand All @@ -115,8 +130,6 @@ leave the video packs unchanged. See [Data Evolution Maintenance](./data-evoluti
`.blob`.
- Non-null writes must be exact descriptor-backed `BlobRef` values containing a
`VideoFrameDescriptor`. Inline bytes and ordinary `BlobDescriptor` values are rejected.
- Version 1 addresses frames by zero-based presentation-order ordinal with stride one. It does not
store PTS values or parse codec/container metadata.
- A `BlobConsumer` callback is not supported for the video field.

## Read Frames
Expand Down
11 changes: 5 additions & 6 deletions docs/docs/pypaimon/lerobot.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,9 +107,8 @@ together, and keep writers paused until the call returns.
Scalars map to scalar types, vectors to `VECTOR`, higher-rank tensors to nested
`ARRAY`, and images to `BLOB`. Images keep their compressed bytes.

Video features map to `BLOB`. Frame rows reference MP4 payloads copied once per
aligned file group. Video imports use the video grouping policy and check
rolling before each Episode. They require a bucket-unaware table. Use
Video features map to `BLOB`. Each MP4 is copied once per aligned file group and indexed for range
reads. Imports require a bucket-unaware table and check rolling before each Episode. Use
`VideoFrameCollator` for scans or `PaimonLeRobotDataset` for training.

## Capture LeRobot frames directly into Paimon
Expand Down Expand Up @@ -223,9 +222,9 @@ loader = DataLoader(dataset, batch_size=32, shuffle=True, num_workers=4)
```

Without `tag_name`, the latest snapshots are used. Frame lookups use the BTree
on `index`; payloads remain lazy. Video decoding prefers TorchCodec, falls back
to PyAV, and reuses a bounded decoder cache. Set `video_backend` to force
either decoder.
on `index`; payloads remain lazy. Indexed videos prefetch metadata and target GOPs, then fetch
uncached PyAV reads on demand. Unindexed videos use the TorchCodec/PyAV scan path. Set
`video_backend` to force either decoder.

Subclass `PaimonDatasetReader` for a custom logical frame layout:

Expand Down
3 changes: 3 additions & 0 deletions docs/docs/pypaimon/video.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ video = pm.BlobDescriptor(
frames.add_video(video, episode_43_rows)
```

With PyAV, PyPaimon indexes supported MP4 files for range reads; other videos
use the scan fallback.

The writer deduplicates exact payload descriptor identity inside each `.video`
file. Its video grouping policy coordinates normal, BLOB, and vector rolling
at payload boundaries. A file may exceed its target before the next boundary.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,15 +37,31 @@ public class VideoFrameDescriptor extends BlobDescriptor {
private static final long MAGIC = 0x564944454F46524DL; // "VIDEOFRM"
private static final byte CURRENT_VERSION = 1;
private static final int FIXED_LENGTH =
Byte.BYTES + Long.BYTES + Integer.BYTES + 3 * Long.BYTES;
Byte.BYTES + Long.BYTES + Integer.BYTES + 5 * Long.BYTES;

private final long frameIndex;

public VideoFrameDescriptor(String uri, long offset, long length, long frameIndex) {
private final long keyframeIndexOffset;
private final long keyframeIndexLength;

public VideoFrameDescriptor(
String uri,
long offset,
long length,
long frameIndex,
long keyframeIndexOffset,
long keyframeIndexLength) {
super(uri, offset, length);
checkArgument(
frameIndex >= 0, "Video frame index must be non-negative, but was %s.", frameIndex);
checkArgument(
keyframeIndexLength >= 0, "Video keyframe index length must be non-negative.");
checkArgument(
(keyframeIndexLength == 0 && keyframeIndexOffset == -1)
|| (keyframeIndexLength > 0 && keyframeIndexOffset >= 0),
"Invalid video keyframe index range.");
this.frameIndex = frameIndex;
this.keyframeIndexOffset = keyframeIndexOffset;
this.keyframeIndexLength = keyframeIndexLength;
}

public long frameIndex() {
Expand All @@ -57,6 +73,22 @@ public BlobDescriptor payloadDescriptor() {
return new BlobDescriptor(uri(), offset(), length());
}

public @Nullable BlobDescriptor keyframeIndexDescriptor() {
return keyframeIndexLength == 0
? null
: new BlobDescriptor(uri(), keyframeIndexOffset, keyframeIndexLength);
}

/** Returns the persisted keyframe index carried by an exact frame reference. */
public static @Nullable Blob keyframeIndexBlob(Blob blob) {
VideoFrameDescriptor frame = fromBlob(blob);
BlobDescriptor mapping = frame == null ? null : frame.keyframeIndexDescriptor();
if (mapping == null) {
return null;
}
return Blob.fromDescriptor(((BlobRef) blob).uriReader(), mapping);
}

/** Returns the video frame carried by an exact lazy blob reference, or {@code null}. */
public static @Nullable VideoFrameDescriptor fromBlob(@Nullable Blob blob) {
if (blob == null || blob.getClass() != BlobRef.class) {
Expand Down Expand Up @@ -86,6 +118,8 @@ public byte[] serialize() {
buffer.putLong(offset());
buffer.putLong(length());
buffer.putLong(frameIndex);
buffer.putLong(keyframeIndexOffset);
buffer.putLong(keyframeIndexLength);
return buffer.array();
}

Expand All @@ -109,15 +143,15 @@ public static VideoFrameDescriptor deserialize(byte[] bytes) {
throw invalidPayload("missing magic header");
}
int uriLength = buffer.getInt();
// checked by comparison and subtraction: uriLength + 3 * Long.BYTES wraps negative
// checked by comparison and subtraction: uriLength + 5 * Long.BYTES wraps negative
// for a uriLength near Integer.MAX_VALUE
if (uriLength < 0) {
throw invalidPayload("negative URI length: " + uriLength);
}
if (uriLength > buffer.remaining()) {
throw invalidPayload("URI length exceeds data size");
}
if (buffer.remaining() - uriLength < 3 * Long.BYTES) {
if (buffer.remaining() - uriLength < 5 * Long.BYTES) {
throw invalidPayload("missing offset/length/frame index");
}

Expand All @@ -127,13 +161,16 @@ public static VideoFrameDescriptor deserialize(byte[] bytes) {
long offset = buffer.getLong();
long length = buffer.getLong();
long frameIndex = buffer.getLong();
long keyframeIndexOffset = buffer.getLong();
long keyframeIndexLength = buffer.getLong();
if (buffer.hasRemaining()) {
throw invalidPayload("trailing bytes");
}
if (frameIndex < 0) {
throw invalidPayload("negative frame index: " + frameIndex);
}
return new VideoFrameDescriptor(uri, offset, length, frameIndex);
return new VideoFrameDescriptor(
uri, offset, length, frameIndex, keyframeIndexOffset, keyframeIndexLength);
}

public static boolean isVideoFrameDescriptor(byte[] bytes) {
Expand All @@ -154,12 +191,13 @@ public boolean equals(Object o) {
}
VideoFrameDescriptor that = (VideoFrameDescriptor) o;
return frameIndex == that.frameIndex
&& payloadDescriptor().equals(that.payloadDescriptor());
&& payloadDescriptor().equals(that.payloadDescriptor())
&& Objects.equals(keyframeIndexDescriptor(), that.keyframeIndexDescriptor());
}

@Override
public int hashCode() {
return Objects.hash(payloadDescriptor(), frameIndex);
return Objects.hash(payloadDescriptor(), frameIndex, keyframeIndexDescriptor());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ public class VideoFrameDescriptorTest {
@Test
public void testRoundTripAndPayloadIdentity() {
VideoFrameDescriptor frame =
new VideoFrameDescriptor("oss://bucket/source.mp4", 17, 103, 42);
new VideoFrameDescriptor("oss://bucket/source.mp4", 17, 103, 42, -1, 0);

assertThat(VideoFrameDescriptor.isVideoFrameDescriptor(frame.serialize())).isTrue();
assertThat(BlobDescriptor.isBlobDescriptor(frame.serialize())).isFalse();
Expand All @@ -46,34 +46,30 @@ public void testRoundTripAndPayloadIdentity() {
.isEqualTo(new BlobDescriptor("oss://bucket/source.mp4", 17, 103));

VideoFrameDescriptor next =
new VideoFrameDescriptor("oss://bucket/source.mp4", 17, 103, 43);
new VideoFrameDescriptor("oss://bucket/source.mp4", 17, 103, 43, -1, 0);
assertThat(next).isNotEqualTo(frame);
assertThat(next.payloadDescriptor()).isEqualTo(frame.payloadDescriptor());

VideoFrameDescriptor indexed =
new VideoFrameDescriptor("oss://bucket/source.mp4", 17, 103, 42, 120, 8);
assertThat(VideoFrameDescriptor.deserialize(indexed.serialize())).isEqualTo(indexed);
assertThat(indexed.keyframeIndexDescriptor())
.isEqualTo(new BlobDescriptor("oss://bucket/source.mp4", 120, 8));
}

@Test
public void testCrossLanguageWireFixture() throws Exception {
VideoFrameDescriptor expected = new VideoFrameDescriptor("s3://bucket/视频.mp4", 7, 99, 42);
byte[] fixture =
fromHex(
new String(
IOUtils.readFully(
VideoFrameDescriptorTest.class
.getClassLoader()
.getResourceAsStream(
"org/apache/paimon/data/video-frame-descriptor-v1.hex"),
true),
StandardCharsets.UTF_8)
.trim());

VideoFrameDescriptor expected =
new VideoFrameDescriptor("s3://bucket/视频.mp4", 7, 99, 42, 106, 8);
byte[] fixture = fixture("video-frame-descriptor-v1.hex");
assertThat(expected.serialize()).isEqualTo(fixture);
assertThat(BlobDescriptor.deserialize(fixture)).isEqualTo(expected);
assertThat(BlobDescriptor.isSerializedDescriptor(fixture)).isTrue();
}

@Test
public void testBlobFromBytesPreservesFrameDescriptor() {
VideoFrameDescriptor expected = new VideoFrameDescriptor("file:/video.mp4", 0, 9, 7);
VideoFrameDescriptor expected = new VideoFrameDescriptor("file:/video.mp4", 0, 9, 7, -1, 0);
Blob blob = Blob.fromBytes(expected.serialize(), null, null);

assertThat(blob).isInstanceOf(BlobRef.class);
Expand All @@ -82,17 +78,18 @@ public void testBlobFromBytesPreservesFrameDescriptor() {

@Test
public void testRejectInvalidPayload() {
VideoFrameDescriptor descriptor = new VideoFrameDescriptor("file:/video.mp4", 0, 9, 7);
VideoFrameDescriptor descriptor =
new VideoFrameDescriptor("file:/video.mp4", 0, 9, 7, -1, 0);
byte[] trailing = Arrays.copyOf(descriptor.serialize(), descriptor.serialize().length + 1);

assertThatThrownBy(() -> VideoFrameDescriptor.deserialize(trailing))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("trailing bytes");
assertThatThrownBy(() -> new VideoFrameDescriptor("file:/video.mp4", 0, 9, -1))
assertThatThrownBy(() -> new VideoFrameDescriptor("file:/video.mp4", 0, 9, -1, -1, 0))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("non-negative");

// The old check let this through because uriLength + 3 * Long.BYTES wrapped
// The old check let this through because uriLength + 5 * Long.BYTES wrapped
// negative, and deserialize went on to allocate ~2GB.
byte[] hostileUriLength = descriptor.serialize();
ByteBuffer.wrap(hostileUriLength)
Expand Down Expand Up @@ -129,4 +126,17 @@ private static byte[] fromHex(String hex) {
}
return bytes;
}

private static byte[] fixture(String name) throws Exception {
return fromHex(
new String(
IOUtils.readFully(
VideoFrameDescriptorTest.class
.getClassLoader()
.getResourceAsStream(
"org/apache/paimon/data/" + name),
true),
StandardCharsets.UTF_8)
.trim());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,4 +15,4 @@
# specific language governing permissions and limitations
# under the License.

014d52464f454449561600000073333a2f2f6275636b65742fe8a786e9a2912e6d7034070000000000000063000000000000002a00000000000000
014d52464f454449561600000073333a2f2f6275636b65742fe8a786e9a2912e6d7034070000000000000063000000000000002a000000000000006a000000000000000800000000000000
Loading
Loading