多模态数据处理快速入门

更新时间:
复制 MD 格式

本文以具身智能的训练数据处理为例,带您在 EMR Serverless Spark 的 Notebook 中用 Daft 跑通一条完整的多模态数据处理链路:从 OSS 读取第一视角视频、分布式抽帧、去除画面黑边,再调用多模态大模型识别每一帧的手部交互动作,最后在 Notebook 中查看图文对照的结果。

背景信息

训练一个能理解厨房操作的具身智能模型,通常需要采集数千小时由人类佩戴摄像头拍摄的第一视角操作视频。这些视频解码后会膨胀成数亿张图像帧,其中还混杂着大量静止画面、模糊镜头和无效空帧。如何从这些冗余数据中高效筛出高质量训练样本,是这类场景的主要瓶颈。

用传统方式做这件事,需要自己拼装多套系统:用 FFmpeg 抽帧转码、写脚本做画面清洗、再对接模型服务做内容打标。本文改用 EMR Serverless Spark 的多模态算子市场,在一条 Daft DataFrame 链路里完成抽帧、清洗和打标,无需自行部署模型服务。

本文用到的能力分为三类:

  • Daft 原生 APIdaft.read_video_frames() 负责分布式视频抽帧。

  • 多模态算子ImageBlackBorderCrop 负责检测并裁掉画面黑边。

  • 表达式快捷 APIai_query() 负责调用多模态大模型生成动作描述。模型服务由 EMR Serverless Spark 内置的 Model Manager 提供,详情请参见模型服务

前提条件

开始前请确认以下条件均已满足:

条件

说明

已创建工作空间

操作方法请参见创建工作空间

已开通 AI 中心

在工作空间左侧导航栏单击 AI中心,然后单击 开通AI中心。相关说明请参见AI概述

已创建 OSS Bucket

用于存放输入视频和输出的图片帧。Bucket 必须与工作空间在同一地域,否则作业无法访问数据。

工作空间按量配额上限不低于 20 CU

配额不足时抽帧和模型调用无法并发执行。调整方法请参见管理工作空间

队列弹性并发上限不低于 20 CU

调整方法请参见管理资源队列

步骤一:上传视频到 OSS

准备一段第一视角拍摄的操作视频(MP4 格式),上传到与工作空间同地域的 OSS Bucket。同时在 Bucket 中规划一个输出目录,用于存放抽帧和裁剪后的图片。如果暂时没有合适的素材,可以直接使用本文的示例视频P01_101_part012_299.40s_310.00s.mp4,这是一段 10 秒左右的第一视角厨房操作片段。

例如:

  • 输入视频:oss://<YOUR_BUCKET>/dataset/kitchen-demo.mp4

  • 输出目录:oss://<YOUR_BUCKET>/dataset/daft-output/

说明

首次体验时建议使用一段 10 秒左右的短视频,可以更快看到完整结果。视频越长、抽帧间隔越小,产生的图片帧和模型调用次数越多,作业耗时也越长。

步骤二:创建 Notebook 会话

登录 EMR Serverless Spark 控制台,进入目标工作空间,创建一个 Notebook 会话并启动。创建时需要注意引擎版本:

  • 引擎版本选择 esr-4.9.0(Spark 3.5.2, Scala 2.12, Ray 2.55.1, Daft 0.7.17)。多模态算子随引擎镜像提供,版本不匹配时会提示算子不存在。

  • 会话的资源规格建议不低于前提条件中的配额要求,否则抽帧与模型调用会排队。

创建 Notebook 会话和运行 Notebook 的详细操作请参见Notebook开发快速入门。引擎版本的管理方式请参见管理运行环境

步骤三:运行数据处理链路

在 Notebook 中新建一个 Python 单元格,粘贴以下代码,并把 INPUT_VIDEOOUTPUT_DIR 替换为您在步骤一中准备的实际路径,然后运行。您也可以下载完整的 Notebook 文件 daft_demo.ipynb,导入 Notebook 会话后直接运行,无需逐段粘贴代码。

import base64
import os
from pathlib import PurePosixPath

import daft
from daft import DataType, col, lit
from daft.emr.functions import ImageBlackBorderCrop, ai_query, emr_udf
from pypaimon.daft import read_paimon, write_paimon

# ==========================================
# 1. 全局配置与 Prompt
# ==========================================
# 请替换为您的 OSS 输入视频和输出目录
INPUT_VIDEO = "oss://<YOUR_BUCKET>/dataset/kitchen-demo.mp4"
OUTPUT_DIR = "oss://<YOUR_BUCKET>/dataset/daft-output/"

PROMPT = """
Role: 你是一位分析第一视角人类视频(Egocentric human videos)的专家。请识别出视频中人的手臂与物体交互的每一个独立动作片段. 
输出要求:输出双手操作的动作片段,英文,20字以内。 
"""


# ==========================================
# 2. 自定义 UDF 定义
# ==========================================
@daft.func(return_dtype=DataType.string())
def image_name(path: str, frame_index: int) -> str:
    """按视频名组织输出帧,避免不同视频的帧混在一起"""
    return f"{PurePosixPath(path).stem}/frame_{frame_index:06d}"


@daft.func(return_dtype=DataType.binary())
def decode_image(value: str) -> bytes:
    """将裁剪结果(base64)转成 Paimon BLOB(bytes)"""
    return base64.b64decode(value)


# ==========================================
# 3. 核心数据处理链路 (抽帧 → 去黑边 → 大模型打标)
# ==========================================
results = (
    # Step 1: 视频抽帧
    daft.read_video_frames(
        path=INPUT_VIDEO,
        image_height=580,
        image_width=740,
        # start_time=600.0,  # 时间范围下推,从第 10 分钟开始读取
        sample_interval_seconds=1,  # 每 2 秒抽取一帧
        decode_thread_type="AUTO",
        decode_thread_count=4,  # 多线程视频解码
    )
    # Step 2: 去黑边 (调用 EMR UDF)
    .with_column(
        "crop",
        emr_udf(
            ImageBlackBorderCrop,
            construct_args={
                "image_src_type": "image_ndarray",
                "detect_algorithm": "auto",
                "crop_sides": ["top", "bottom"],  # 支持选择裁切黑边的方向
                "output_dir": OUTPUT_DIR,
                "image_format": "jpg",
                "quality": 85,
            },
            concurrency=2,
            batch_size=32,
        )(
            image=col("data"),
            image_name=image_name(col("path"), col("frame_index")),
        ),
    )
    # Step 3: 提取图片 URL 与 Base64 数据
    .with_column("image_url", col("crop").get("image_path"))
    .with_column("frame_blob", decode_image(col("crop").get("base64")))
    # Step 4: 多模态大模型直接读取 OSS 图片并生成场景描述
    .with_column(
        "scene_description",
        ai_query(lit(PROMPT), data=col("image_url"), data_type="uri", service_name="qwen3.7-plus").get("content"),
    )
    # Step 5: 字段筛选与重命名 (只保留最终需要写入 Paimon 的字段)
    .select(
        col("path").alias("video_path"),
        col("frame_index"),
        col("frame_time"),
        col("image_url"),
        col("frame_blob"),
        col("scene_description"),
        lit("2026-08-01").alias("dt"),
    )
)

# ==========================================
# 4. 触发执行与结果打印
# ==========================================
results.collect()
print(results)

这段代码的处理链路分为四步:

步骤

使用的能力

说明

视频抽帧

daft.read_video_frames()

从 OSS 读取视频并按 sample_interval_seconds 指定的间隔抽帧,同时统一输出分辨率。取消注释 start_time 可以只处理视频的后半段。

去黑边

ImageBlackBorderCrop 算子

自动检测并裁掉画面上下的黑边,裁剪后的图片按视频名分目录写入 OUTPUT_DIR。算子通过 emr_udf() 包装后调用。

大模型打标

ai_query()

把图片的 OSS 路径以 data_type="uri" 的方式传给 qwen3.7-plus,由模型按 PROMPT 输出每帧的手部动作描述。

字段整理

select()

只保留视频路径、帧序号、帧时间、图片地址、图片二进制和动作描述六个字段,便于后续写入数据湖表。

代码中还定义了两个自定义 UDF:image_name() 按视频名组织输出帧的命名,避免多个视频的帧混在同一目录;decode_image() 把算子返回的 Base64 字符串还原为二进制,供后续展示或写表使用。

步骤四:查看处理结果

新建单元格,运行以下代码查看 DataFrame 的结构与内容:

results

再新建一个单元格,运行以下代码在 Notebook 中以图文对照的方式展示前 6 帧的画面和对应的动作描述:

from emrssutils.utils import show_df_image

show_df_image(results, image_row_name='frame_blob', time_row_name='frame_time', description_row_name='scene_description', limit=6)

输出中每一行对应一个视频帧:左侧是去黑边后的画面,右侧是模型按 PROMPT 生成的英文动作描述。裁剪后的图片同时保留在 OUTPUT_DIR 目录中,可直接作为下游训练样本。

后续步骤

  • 把 Notebook 中调试好的链路改为批任务或提交到 Ray 集群,用于处理更大规模的数据,请参见Ray集群提交任务