Skip to content
Merged
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
41 changes: 23 additions & 18 deletions docs/docs/concepts/spec/fileformat.md
Original file line number Diff line number Diff line change
Expand Up @@ -401,9 +401,8 @@ Limitations:

### Video

Video is an independent, versioned format with the `.video` extension. It packs one or more
complete encoded-video payloads and logical frame runs. The payloads are raw byte ranges without
the ordinary BLOB entry header, length trailer, or per-entry CRC:
`.video` stores complete encoded videos and frame runs, without BLOB entry headers, length
trailers, or per-entry CRC. On-disk order:

```
+----------------------------+
Expand All @@ -413,15 +412,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 block lengths per video (0 = scan fallback)
+----------------------------+
| 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 @@ -430,14 +436,18 @@ the ordinary BLOB entry header, length trailer, or per-entry CRC:
+----------------------------+
```

The run arrays have equal element counts. A non-negative run reference is an ordinal in the
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.
Run arrays have equal lengths. Non-negative references select a video; `-1` means NULL and `-2`
means a data-evolution placeholder. Row `r` in a run starting at `s` maps to frame
`run_first_frame + r - s`; gaps start new runs.

The serialized `VideoFrameDescriptor` stored in an Arrow/data-file cell has its own versioned
wire layout. All numeric values are little endian:
A keyframe-index block contains a 17-byte header (version `1`: uint8; magic `0x564944454F4B4649`:
uint64; metadata-range and keyframe counts: uint32 each), metadata `(offset, length)` pairs
(int64 each), and zlib-compressed `(frame ordinal, PTS, packet position)` entries (int64 each).
Numeric fields are little endian; offsets are payload-relative. Limits: 65,536 metadata ranges
and 65,536 keyframes, 16 MiB per block, 64 MiB per file. The index covers the first video
stream; its time base stays in the video.

An Arrow/data-file cell stores a separately versioned, little-endian `VideoFrameDescriptor`:

| Field | Size | Description |
| --- | ---: | --- |
Expand All @@ -448,14 +458,9 @@ 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.

Readers validate footer and index bounds, positive physical lengths, full coverage of the payload
region, equal run-index counts, positive run lengths, physical ordinals, and non-negative first
frames. The format currently supports one scalar BLOB field per file. Physical video reuse uses
exact input payload `BlobDescriptor` identity and is file-local; there are no cross-file payload
references. Ordinary `.blob` files keep their existing version, wrappers, checksums, and layout.
A `.video` file serves one scalar BLOB field; references are file-local. `.blob` is unchanged.

For usage details, configuration options, and examples, see [Blob Type](../../multimodal-table/blob).
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 @@ -115,6 +115,9 @@ video = pm.BlobDescriptor(
frames.add_video(video, episode_43_rows)
```

With PyAV, PyPaimon builds keyframe indexes for supported ISO BMFF videos (such
as MP4); other decodable 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());
}
}
Loading
Loading