本文介绍 Python DataFrame API 中的视频元数据、视频切分和视频抽帧算子。
使用限制
仅实时计算引擎 VVR 11.8 及以上版本支持。
视频读取和解码使用 PyAV/FFmpeg。
视频 URI 必须需要能通过 Flink FileSystem 访问。
video_split和video_explode_frames是 UDTF,需要通过DataFrame.join_lateral使用。本文使用 视频片段引用类型 (VIDEO_STRUCT_TYPE)、视频元数据类型 (VIDEO_METADATA_TYPE) 和 视频帧元数据类型 (VIDEO_FRAME_METADATA_TYPE) 表示
pyflink.multimodal.types中定义的 DataFrame API 类型。
算子清单
分类 | 算子 | 说明 |
元数据 | 读取视频基本信息,例如分辨率、帧率、时长和编码格式。该算子只读取容器和视频流元数据,不解码视频帧。 | |
视频切分 | 将完整视频或已有视频片段切分为多个片段引用,展开为一行一个片段引用。该算子本身不复制视频内容,也不读取视频帧,只计算时间窗口,产出引用片段。 | |
视频抽帧 | 从完整视频或视频片段中抽帧,并将结果展开为一帧一行。 | |
从完整视频或视频片段中抽帧,并把同一个输入视频的结果收集到同一行的数组中。若视频较长或抽帧数量较多,推荐使用 |
数据类型说明
视频片段引用类型 (VIDEO_STRUCT_TYPE)
VIDEO_STRUCT_TYPE 的字段结构如下:
DataType.struct({
"uri": DataType.string(),
"start_time_ms": DataType.int64(),
"end_time_ms": DataType.int64()
})字段 | DataFrame API 类型 | 说明 |
|
| 视频文件路径或对象存储 URI。 |
|
| 片段开始时间,单位为毫秒,包含该时间点。 |
|
| 片段结束时间,单位为毫秒,包含该时间点。 |
片段引用只描述闭区间 [start_time_ms, end_time_ms],不复制视频内容。
视频元数据类型 (VIDEO_METADATA_TYPE)
VIDEO_METADATA_TYPE 的字段结构如下:
DataType.struct({
"width": DataType.int32(),
"height": DataType.int32(),
"fps": DataType.float64(),
"duration_ms": DataType.int64(),
"frame_count": DataType.int64(),
"time_base": DataType.float64(),
"codec_name": DataType.string(),
"video_stream_index": DataType.int32()
})字段 | DataFrame API 类型 | 说明 |
|
| 视频宽度,单位为像素。 |
|
| 视频高度,单位为像素。 |
|
| 视频帧率。 |
|
| 视频时长,单位为毫秒。 |
|
| 视频帧数。 |
|
| 视频流时间基准,用于计算帧时间。 |
|
| 视频编码格式名称,例如 |
|
| 选中的视频流编号。 |
视频帧元数据类型 (VIDEO_FRAME_METADATA_TYPE)
VIDEO_FRAME_METADATA_TYPE 的字段结构如下:
DataType.struct({
"uri": DataType.string(),
"video_stream_index": DataType.int32(),
"frame_index": DataType.int64(),
"pts": DataType.int64(),
"time_ms": DataType.int64(),
"key_frame": DataType.boolean(),
"start_time_ms": DataType.int64(),
"end_time_ms": DataType.int64(),
})字段 | DataFrame API 类型 | 说明 |
|
| 原始视频文件路径或对象存储 URI。 |
|
| 输出帧所属的视频流编号。 |
|
| 输出帧在当前抽帧结果中的序号,从 0 开始。 |
|
| 视频文件中的原始帧时间戳,通常用于和视频处理工具对齐;无法确定时为 |
|
| 帧对应的视频时间,单位为毫秒;无法确定时为 |
|
| 是否为关键帧。 |
|
| 当前抽帧输入片段的开始时间,单位为毫秒;输入为完整 URI 时为 |
|
| 当前抽帧输入片段的结束时间,单位为毫秒;输入为完整 URI 时为 |
通用 Runtime 参数
函数签名保留各算子实际支持的 Runtime 参数。为避免重复,算子参数表只说明业务参数,Runtime 参数统一说明如下。
参数 | 类型 | 默认值 | 适用范围 | 说明 |
|
|
| 本页全部算子 | UDF 或 UDTF 并发度。 |
视频元数据与切分
video_metadata
读取视频基本信息,例如分辨率、帧率、时长和编码格式。该算子只读取容器和视频流元数据,不解码视频帧。
输入类型: DataType.string(),视频 URI 列。
函数签名:
video_metadata(
*columns,
on_error="raise",
container_options=None,
read_chunk_size=None,
max_cached_blocks=None,
read_ahead_blocks=None,
concurrency=None
)参数 | 类型 | 默认值 | 说明 |
|
|
|
|
|
|
| 传给 PyAV |
|
|
| Java FileSystem Bridge 的单次读取字节数,必须大于 0 且不能超过 16 MiB。 |
|
|
| 每个文件最多缓存的块数; |
|
|
| 缓存未命中后预读的块数; |
空值 URI 返回 None。
返回类型: VIDEO_METADATA_TYPE。
from pyflink.dataframe import col
from pyflink.multimodal.operators import video_metadata
result = df.with_column(
"video_metadata",
video_metadata(
col("uri"),
on_error="null"
)
)video_split
将完整视频或已有视频片段切分为多个片段引用,展开为一行一个片段引用。该算子本身不复制视频内容,也不读取视频帧,只计算时间窗口。
输入类型: DataType.string() 或 VIDEO_STRUCT_TYPE;可追加 VIDEO_METADATA_TYPE 或 DataType.int64(),第一个输入是视频 URI 或已有片段引用。第二个输入可提供视频元数据或毫秒时长,避免重复探测。空值视频输入不输出行。
函数签名:
video_split(
*columns,
segment_duration_ms=None,
num_segments=None,
video_duration_ms=None,
max_segments=1024,
on_error="raise",
container_options=None,
read_chunk_size=None,
max_cached_blocks=None,
read_ahead_blocks=None,
concurrency=None
)参数 | 类型 | 默认值 | 说明 |
|
|
| 固定片段长度,单位为毫秒,必须大于 0。与 |
|
|
| 近似等长的目标片段数,必须大于 0。与 |
|
|
| 显式视频时长,单位为毫秒;在输入没有结束时间且第二个输入没有有效时长时作为后备值。 |
|
|
| 单行最多输出的片段数,必须大于 0。达到限制后停止输出剩余片段。 |
|
|
|
|
|
|
| 算子自行探测元数据时传给 PyAV |
|
|
| 算子自行探测元数据时的单次读取字节数,不能超过 16 MiB。 |
|
|
| 每个文件最多缓存的块数; |
|
|
| 缓存未命中后预读的块数; |
按固定时长切分时,算子生成闭区间片段。例如,2 秒视频按 1000 毫秒切分为 [0, 999] 和 [1000, 1999]。按目标片段数切分时使用向上取整的片段时长,因此实际片段数可能少于目标值。
返回类型: UDTF 单列 segment,类型为 VIDEO_STRUCT_TYPE。
from pyflink.multimodal.operators import video_split
segments = df.join_lateral(
video_split(
col("uri"),
segment_duration_ms=10_000
).alias("segment")
)如果已有 video_metadata 的结果,可以作为第二个输入列传入:
with_metadata = df.with_column(
"metadata",
video_metadata(col("uri")),
)
segments = with_metadata.join_lateral(
video_split(
col("uri"),
col("metadata"),
segment_duration_ms=10_000
).alias("segment")
)视频抽帧
video_explode_frames
从完整视频或视频片段中抽帧,并将结果展开为一帧一行。
输入类型: DataType.string() 或 VIDEO_STRUCT_TYPE,视频 URI 列或 video_split 返回的视频片段引用列。空值输入不输出行。
函数签名:
video_explode_frames(
*columns,
frame_selector="all_frames",
sample_interval_ms=None,
max_frames=None,
image_height=None,
image_width=None,
on_error="raise",
container_options=None,
read_chunk_size=None,
max_cached_blocks=None,
read_ahead_blocks=None,
concurrency=None
)参数 | 类型 | 默认值 | 说明 |
|
|
| 抽帧策略。支持 |
|
|
|
|
|
|
| 每个输入最多输出的帧数,必须大于 0; |
|
|
| 输出帧高度,单位为像素,必须与 |
|
|
| 输出帧宽度,单位为像素,必须与 |
|
|
|
|
|
|
| 传给 PyAV |
|
|
| Java FileSystem Bridge 的单次读取字节数,不能超过 16 MiB。 |
|
|
| 每个文件最多缓存的块数; |
|
|
| 缓存未命中后预读的块数; |
输入为视频片段引用时,算子按闭区间 [start_time_ms, end_time_ms] 选择时间已知的帧。"sample" 从片段起点建立采样时间点,并选择时间不早于各采样点的第一帧。
返回类型: UDTF 两列:frame 为 DataType.image()(帧图像),frame_metadata 为 VIDEO_FRAME_METADATA_TYPE。
from pyflink.multimodal.operators import video_explode_frames
frames = df.join_lateral(
video_explode_frames(
col("uri"),
frame_selector="sample",
sample_interval_ms=1_000,
max_frames=300,
image_height=360,
image_width=640,
on_error="skip"
).alias("frame", "frame_metadata")
)video_extract_frames
从完整视频或视频片段中抽帧,并把同一个输入视频的结果收集到同一行的数组中。若视频较长或抽帧数量较多,推荐使用 video_explode_frames,避免单行持有过多图像。
输入类型: DataType.string()(视频 URI 列)或 VIDEO_STRUCT_TYPE(video_split 返回的视频片段引用列)。
函数签名:
video_extract_frames(
*columns,
frame_selector="all_frames",
sample_interval_ms=None,
max_frames=None,
image_height=None,
image_width=None,
on_error="raise",
container_options=None,
read_chunk_size=None,
max_cached_blocks=None,
read_ahead_blocks=None,
concurrency=None
)参数 | 类型 | 默认值 | 说明 |
|
|
| 抽帧策略。支持 |
|
|
|
|
|
|
| 保留在同一输出行中的帧数上限,必须大于 0; |
|
|
| 输出图像高度,必须与 |
|
|
| 输出图像宽度,必须与 |
|
|
|
|
|
|
| 传给 PyAV |
|
|
| Java FileSystem Bridge 的单次读取字节数,不能超过 16 MiB。 |
|
|
| 每个文件最多缓存的块数; |
|
|
| 缓存未命中后预读的块数; |
空值输入或 on_error="null" 下的处理失败返回 None;处理成功但没有选中帧时返回空数组。
返回类型: DataType.list(DataType.image())。
from pyflink.multimodal.operators import video_extract_frames
# 1.直接从 URI 读取并抽帧
result = df.with_column(
"sampled_frames",
video_extract_frames(
col("uri"),
frame_selector="sample",
sample_interval_ms=1_000,
max_frames=60
)
)
# 2.先构造 CLIP ref,每个 CLIP 10s
clips = df.join_lateral(
video_split(
col("uri"),
col("metadata"),
segment_duration_ms=10_000,
max_segments=1024,
on_error="skip"
).alias("clip")
)
# 利用 CLIP ref 分段抽帧
result = clips.with_column(
"sampled_frames",
video_extract_frames(
col("clip"),
frame_selector="sample",
sample_interval_ms=1_000,
max_frames=60
)
)完整 Pipeline 示例
以下示例先读取元数据并切成 10 秒片段,再从每个片段中每秒抽取一帧。
from pyflink.dataframe import col
from pyflink.multimodal.operators import (
video_explode_frames,
video_metadata,
video_split
)
with_metadata = df.with_column(
"metadata",
video_metadata(
col("uri"),
on_error="null"
)
).filter(col("metadata").is_not_null)
segments = with_metadata.join_lateral(
video_split(
col("uri"),
col("metadata"),
segment_duration_ms=10_000,
max_segments=1024,
on_error="skip"
).alias("segment")
)
frames = segments.join_lateral(
video_explode_frames(
col("segment"),
frame_selector="sample",
sample_interval_ms=1_000,
max_frames=10,
on_error="skip"
).alias("frame", "frame_metadata")
)