← 返回AI教程
🌐 其他

Daft 多模态视频抽帧性能优化实践:从抽帧到大模型打标,一条视频理解流水线是怎么跑起来的

来源:掘金 · 发布于 2026-08-18 18:13:03
EMR Serverless Spark 团队推出了深度优化的 Daft 版本,不仅大幅提升了视频、图片等多模态数据的读取性能,还提供了 ai_query 等大模型调用原语和丰富的内置算子。

Daft 多模态视频抽帧性能优化实践:从抽帧到大模型打标,一条视频理解流水线是怎么跑起来的

阿里云大数据AI技术 2026-08-18 0 阅读11分钟

前言

过去两年,多模态大模型的能力跃升让"把视频喂给模型,让它看懂画面里发生了什么"成为越来越多用户的实际需求;与此同时,监控录像、直播回放、UGC 内容、工业质检影像等场景产生的视频数据动辄以 TB 计。要把这些数据变成可检索、可分析的结构化信息,"视频抽帧 + 大模型理解"成了最直接的路线。

Daft 是一个面向多模态数据的开源分布式 DataFrame 引擎,原生支持图片、视频等非结构化数据的读取和变换。在此基础上,EMR Serverless Spark 团队推出了深度优化的 Daft 版本,不仅大幅提升了视频、图片等多模态数据的读取性能,还提供了 ai_query 等大模型调用原语和丰富的内置算子(覆盖文本清洗、图像处理、格式转换等),让"数据读取 → 预处理 → AI 理解"在一条 DataFrame 流水线里闭环。我们基于这个版本构建了一条从视频到结构化标签的完整管线,并在实际生产中把它打磨到了可用的性能水平。这篇文章会顺着数据的流动,把这条管线从头到尾走一遍。 image - 2026-08-18T181031.926.png

视频理解,正在成为数据管线的标配

内容审核要判断画面是否合规,视频打标要给每一帧生成描述和标签,直播质检要抓取异常片段,安防回看要检索特定场景——这些需求背后都是同一条流水线:把视频拆成一帧帧图片,交给大模型理解,再把结果结构化入库。

先看一个完整的例子。假设 OSS 的某个目录下有一批视频,每个都是以 GB 计的大文件,我们只关心每个视频第 10 分钟之后的内容,希望每 10 秒抽一帧,把帧图片存到 OSS,再让大模型给每帧生成描述标签:

from pathlib import PurePosixPath

import daft
from daft import DataType, col, lit
from daft.emr.functions import ImageBlackBorderCrop, ai_query, emr_udf
from daft.io import IOConfig, JindoConfig

io_config = IOConfig(
    jindo=JindoConfig(
        endpoint="oss-<your-region>-internal.aliyuncs.com",
        access_key_id="<your-ak>",
        access_key_secret="<your-sk>",
    )
)


@daft.func(return_dtype=DataType.string())
def make_image_name(path: str, frame_index: int) -> str:
    """按 <视频名>/frame_<帧号> 命名,避免不同视频的帧混在一起。"""
    video_name = PurePosixPath(str(path).split("?")[0]).stem
    return f"{video_name}/frame_{int(frame_index):06d}"


frames_df = (
    daft.read_video_frames(
        path="oss://my-bucket/videos/*.mp4",   # 目录下所有视频
        image_height=480,
        image_width=640,
        start_time=600.0,               # 只从第 10 分钟开始读
        sample_interval_seconds=10.0,   # 每 10 秒取一帧
        decode_thread_type="AUTO",      # 多线程解码
        decode_thread_count=4,
        io_config=io_config,
    )
    # 去除黑边并写入 OSS(返回 struct{base64, image_path},取 image_path 作为 OSS 路径)
    .with_column(
        "crop_result",
        emr_udf(
            ImageBlackBorderCrop,
            construct_args={
                "image_src_type": "image_ndarray",
                "detect_algorithm": "auto",
                "crop_sides": ["top", "bottom"],
                "output_dir": "oss://my-bucket/frames/",
                "image_format": "jpg",
                "quality": 85,
                "io_config": io_config,
            },
            concurrency=4,
            batch_size=32,
        )(image=col("data"), image_name=make_image_name(col("path"), col("frame_index"))),
    )
    .with_column("image_url", col("crop_result").get("image_path"))
    # 让大模型给每一帧打标
    .with_column(
        "ai_result",
        ai_query(
            lit(
                "请描述这张画面的内容,包括场景布局、人物动作、环境特征、"
                "整体氛围。输出 JSON,包含 scene_description、tags(数组)、"
                "mood 三个字段。"
            ),
            data=col("image_url"),
            data_type="uri",
        ),
    )
    # ai_result 是封装结构体,模型的业务输出在 content 字段里
    .with_column("content", col("ai_result").get("content"))
)

# 小批量验证:collect() 把结果一次性物化出来
results = frames_df.select(
    col("path"), col("frame_time"), col("image_url"), col("content")
).collect()

短短几十行,一条"视频 → 图片 → AI 标注 → 结构化结果"的管线就成型了。下面我们顺着数据的流动,看看每一步发生了什么。

第一步:把视频读进来,但只读需要的部分

流水线的入口是 read_video_frames。它接受一个 OSS 路径,直接从对象存储读取视频,不需要你先把一个大文件下载到本地——远程读取、格式解析、解码这些脏活都由框架接管。

但"读视频"这一步恰恰是整条管线里最容易被低估的瓶颈。视频解码是 CPU 密集型操作,一段 20 分钟、30fps 的视频有 36000 帧,如果从头到尾逐帧顺序解码,光解码就要花费很长时间。

这里藏着我们做的第一类优化:读取时的按需下推。业务往往只关心视频的一部分——上面的例子里我们只要第 10 分钟之后的内容。通过 start_time 把这个信息下推到读取源,Daft 会利用视频容器的关键帧索引直接定位到目标位置附近开始解码,而不是从第 0 秒顺序解到第 10 分钟。前 600 秒的解码工作被整段跳过,对于远程大文件,网络 IO 也不必再从头拉取整个文件前缀。

值得一提的是,这种读取是自动分布式的。上面的例子里,path 用了通配符匹配整个目录,当目录下有几十上百个视频时,Daft 会自动把每个视频的解码、裁边、上传、模型调用拆成独立任务,分配到 Ray 集群的多个 worker 上并行执行;单个视频的帧数据量较大时,还会被进一步切成多个分区。你不需要写任何分布式代码,也不需要关心资源怎么切分、任务怎么调度——这些全部由框架接管。

第二步:采样,只保留有价值的帧

视频读进来之后,业务真正需要的通常不是每一帧,而是按固定节奏采样的关键帧。sample_interval_seconds=10 告诉 Daft:每 10 秒给我一帧就够了。

采样看似只是"每隔几秒挑一帧",但朴素实现里藏着大量浪费。30fps 的视频,10 秒间隔意味着两个采样点之间有 300 帧,其中 299 帧是要被丢掉的。如果顺序解码,这 299 帧全都被完整解出来又立刻扔掉,CPU 白白空转。

这就是我们做的第二类优化:跳过不需要的帧。当采样间隔足够大时,Daft 在输出一帧之后会直接 seek 到下一个采样点附近,跳过中间那几百帧的解码。解码器只碰它真正要输出的那些帧,把"解出来再丢掉"变成"根本不解"。这项优化对可变帧率(VFR)视频同样有效,让手机录像、直播回放这类视频也能受益。

与此同时,剩下那些真正要解码的帧,也不该单线程慢慢啃。decode_thread_type 和 decode_thread_count 让底层解码器用多个线程并行工作,把多核 CPU 用满。这是我们在解码环节做的第三类优化,和采样 seek 配合,让"少解 + 快解"叠加生效。

第三步:画面预处理——去掉黑边,直接落盘

采样出来的帧还不能直接交给大模型。真实场景的视频画面经常带有黑边——录制设备的画幅比例和输出分辨率不匹配时,画面上下或左右就会出现无意义的黑色条带。这些黑边不仅浪费存储空间,还会干扰大模型对画面内容的理解。

EMR Serverless Spark Daft 提供了 ImageBlackBorderCrop 算子来处理这个问题。它能自动检测并裁掉画面四周的黑边,支持多种检测算法(阈值扫描、边缘检测、直方图分析),也可以用 auto 模式让三种算法投票取共识,避免单一算法在特定画面上误判。你还可以通过 crop_sides 指定只处理某几个方向——比如视频只有顶部黑边,就只裁顶部。

更重要的是,这个算子不只是裁剪。它同时承担了编码和上传的工作:裁剪完成后直接把 JPEG 写入你指定的 OSS 路径,并返回一个包含该路径的结构体({base64, image_path})。你不需要自己写"裁剪 → 编码 → 上传"的三步逻辑,一个 emr_udf 调用全部搞定,框架内部自动处理并发和批量。取出其中的 image_path 就是可访问的 OSS 图片路径。

两个调度参数顺带说明:concurrency 指定同时运行的 UDF 实例数(示例中为 4),batch_size 指定每个实例每批处理的行数(示例中为 32),两者共同决定这一步的处理吞吐,可以按 CPU 和网络带宽调整。

这一步之后,每一帧都变成了一张干净的、可访问的 OSS 图片 URL,随时可以交给下游。

第四步:交给大模型打标

图片就位后,ai_query 负责调用大模型。它以 DataFrame 的一列(图片 URL)为输入,你只需要写一段 prompt 描述你想要什么,剩下的批量请求、并发调度、结果解析都由框架处理。

需要注意的是,ai_result 并不是模型的业务输出本身,而是一个封装结构体:除了 content 里的业务内容,还带有 finish_reason、Token 用量、错误信息等运行元数据。我们让模型输出的 scene_description、tags、mood JSON 位于 content 字段(字符串),下游按需解析成结构化字段即可:

# ai_result 是封装结构体,模型的业务输出在 content 字段里
.with_column("content", col("ai_result").get("content"))

到这一步,结果就可以进入下游——继续 join、过滤,或写入 Paimon 等结果表。示例里我们用了 collect() 把结果一次性物化,适合小批量验证;生产环境建议换成写表动作。

整条管线到 collect() 才真正开始执行。在此之前的每一步都只是在构建执行计划,得益于惰性求值,读取、采样、编码、上传、模型调用会被编排成一条连贯的流水线,中间结果不需要全量物化在内存里。

性能实测:三项优化的叠加收益

视频理解管线里,大家的注意力往往都在模型调用端,但如果前面的视频解码还在单线程从头读到尾,再快的模型也只能干等着。我们在读取和采样环节做的优化——按需下推、跳过无用帧、多线程解码——叠加起来收益有多大?我们在同一个 Ray 集群、同一镜像上做了 A/B 对照,两个版本唯一的差异就是优化代码本身:

测试素材是一段存放在 OSS 上、约 1200 秒的 VFR 监控视频(30fps),抽取场景为读取第 900~1200 秒区间、每 1 秒抽一帧:

优化项Baseline (s)优化后 (s)加速比
区间下推13823224.30x
采样跳帧133510871.23x
多线程解码13317891.69x
三项全部叠加133412111.05x

三项优化各自独立、几乎线性叠加:叠加后原本要跑 22 分钟的任务,2 分钟就完成了。

写在最后

性能只是其中一面。回过头看文章中的代码:读视频、去黑边、写 OSS、调大模型,每一步都是一个现成的函数调用,没有一处需要你手写解码循环、编码逻辑、上传重试或并发控制。这才是我们更希望传递的价值——处理数据时,你应该把精力花在业务本身,想清楚"要从数据里得到什么",而不是耗在"怎么把它高效地读进来、拆开、再传出去"这些底层细节上。

EMR Serverless Spark Daft 的高性能算子和 AI Function 正是为此而生。视频抽帧、黑边裁剪、图像处理、文本清洗、大模型调用……这些在数据管线里反复出现、又容易写错写慢的环节,我们把它们沉淀成经过优化的、开箱即用的函数:复杂度留在框架内部,性能、并发、容错由我们兜底,你只需要像搭积木一样把它们串成自己的管线。

去黑边只是冰山一角。EMR Serverless Spark Daft 内置了一个持续扩充的算子市场,覆盖图像、视频、文本、音频等多模态数据的常见处理场景,把"数据读取 → 预处理 → AI 理解"各环节反复出现的脏活累活都沉淀成了开箱即用的函数。下面这张图列举了其中一部分算子:

image.png

可用性说明:本文介绍的 AI Function 目前处于邀测阶段,尚未正式开放。如需抢先体验,请联系 EMR Serverless Spark 团队申请白名单。钉钉群号:58570004119。