本文以具身智能的训练数据处理为例,带您在 EMR Serverless Spark 的 Notebook 中用 Daft 跑通一条完整的多模态数据处理链路:从 OSS 读取第一视角视频、分布式抽帧、去除画面黑边,再调用多模态大模型识别每一帧的手部交互动作,最后在 Notebook 中查看图文对照的结果。
背景信息
训练一个能理解厨房操作的具身智能模型,通常需要采集数千小时由人类佩戴摄像头拍摄的第一视角操作视频。这些视频解码后会膨胀成数亿张图像帧,其中还混杂着大量静止画面、模糊镜头和无效空帧。如何从这些冗余数据中高效筛出高质量训练样本,是这类场景的主要瓶颈。
用传统方式做这件事,需要自己拼装多套系统:用 FFmpeg 抽帧转码、写脚本做画面清洗、再对接模型服务做内容打标。本文改用 EMR Serverless Spark 的多模态算子市场,在一条 Daft DataFrame 链路里完成抽帧、清洗和打标,无需自行部署模型服务。
本文用到的能力分为三类:
Daft 原生 API:
daft.read_video_frames()负责分布式视频抽帧。多模态算子:
ImageBlackBorderCrop负责检测并裁掉画面黑边。表达式快捷 API:
ai_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_VIDEO 和 OUTPUT_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)这段代码的处理链路分为四步:
步骤 | 使用的能力 | 说明 |
视频抽帧 |
| 从 OSS 读取视频并按 |
去黑边 |
| 自动检测并裁掉画面上下的黑边,裁剪后的图片按视频名分目录写入 |
大模型打标 |
| 把图片的 OSS 路径以 |
字段整理 |
| 只保留视频路径、帧序号、帧时间、图片地址、图片二进制和动作描述六个字段,便于后续写入数据湖表。 |
代码中还定义了两个自定义 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集群提交任务。