一批智驾路采视频进入 OSS 后,数据团队真正想做的,往往不是“写一段 Python”,而是尽快得到一份可检索、可标注、可复用的数据集:原始视频按道路场景完成切分,关键帧带上道路环境、天气、交通参与者和异常事件标签,Embedding、处理状态和行级错误写入数据湖。
问题是,从一句业务需求到一条真正跑通的分布式 Pipeline,中间还隔着算子选择、模型调用、OSS 访问、Ray 资源、错误重试、结果落表和远程验收。
现在,这些工程工作可以交给 alibabacloud-emr-daft-multimodal-job Skill。
一句话看懂 向 Qoder 描述智驾数据的输入、处理目标和输出位置,Skill 会根据目标 EMR Serverless Spark 发布版本,生成完整的 EMR Daft 作业包;获得提交授权后,再通过 SubmitRayJob 远程运行并回读结果,完成从 OSS 原始媒体到 Paimon 统一数据资产的处理闭环。
01
一个典型需求可能只有一句话:
“生成一个智驾路采视频数据处理作业:读取 Paimon 视频清单表,表内包含 OSS 视频 URI 和采集元数据;完成场景切分、关键帧提取、道路环境、天气、交通参与者和异常事件标签生成,再为关键帧生成多模态向量;将场景片段、关键帧、标签、多模态向量、处理状态和行级错误分别写入 Paimon 表。”
但落到工程上,至少要回答这些问题:
视频如何批量读取、转码、裁剪、抽帧或切分,避免单机成为瓶颈;
哪些阶段应该放在同一个作业,哪些中间结果需要物化,避免昂贵步骤重复执行;
通用视频理解、专用视频算子和 Embedding 能力应该如何选择;
模型调用如何控制并发、处理限流,并保留行级错误与 token 用量;
OSS 地址如何交给外部模型访问,同时避免把签名 URL 写入结果;
Paimon 表结构与作业输出不一致时,怎样在写入前发现问题;
代码生成后,如何确认它真的能在目标 Workspace 中运行、写入并回读。
这些问题高度重复,却会直接决定智驾数据 Pipeline 能不能稳定交付。
alibabacloud-emr-daft-multimodal-job 把 EMR Daft 的 API、媒体算子、模型服务、分布式运行、存储、安全约束和验收流程沉淀成一套可复用的生成规则,让数据团队把精力放回场景定义、标签体系和数据质量。
02
下面用这条智驾路采视频作业贯穿全文。Skill 不是从“有哪些算子”出发拼功能,而是沿着路采数据的处理链路选择能力:
路采作业的典型处理环节
具体 Pipeline 不会被固定在一套模板里。对于这条路采视频作业,Skill 会先核对目标发布版本实际提供的 API,再从 VideoKeyframeExtract、VideoSceneSeg、VideoUnderstanding、ai_query、ai_embedding_multimodal 等能力中选择与需求匹配的组合。
能力边界 视频理解可以根据业务 Prompt 生成道路环境、天气、交通参与者和异常事件等标签,但标签定义、Prompt 和质量验收标准仍由智驾业务提供。VideoRiskRec 面向视频内容风险与合规审核,不应直接等同于驾驶安全事件识别。
03
对于短而线性的任务,Skill 会保持单作业,减少不必要的中间存储:
当预处理成本较高、结果需要复用,或模型阶段需要独立重跑时,Skill 会拆成多阶段 Pipeline:
这样做有三个直接价值:
可复跑:模型 Prompt 或标签体系变化时,不必重新解码全部原始视频;
可复用:同一份场景片段或关键帧可以服务多个标注、检索和训练任务;
可观测:每个阶段都有明确输入、输出和错误边界,便于抽样检查与失败数据重跑。
04
业务信息
|
运行环境信息
Workspace ID 与地域;
Workspace 支持的 Ray 发布版本;
队列与 head/worker 资源规格;
OSS working directory;
作业主动超时时间;
需要访问 VPC 服务时使用的网络服务。
信息准备好后,可以直接对 Qoder 说:
/alibabacloud-emr-daft-multimodal-job生成一个智驾路采视频数据处理作业:输入是 Paimon 视频清单表,包含 video_uri、vehicle_id 和 collect_time,其中 video_uri 指向 OSS 上的原始视频;先对视频做场景切分和关键帧提取,再识别道路环境、天气、交通参与者和异常事件;保留媒体 URI、标签、模型名称、token 用量、处理状态和行级错误;为关键帧生成多模态 Embedding;将场景片段、关键帧、标签、Embedding 和质量结果分别写入 Paimon 表。使用我提供的 Workspace、发布版本、队列、资源规格和 workingDir。先生成并检查完整作业包;获得我的明确授权后,再通过 SubmitRayJob 提交,等待作业完成并回读各张 Paimon 表的结果。
如果目标发布镜像缺少某个算子,Skill 会先给出版本差异和替代 Pipeline,而不是生成一份只能在本地源码中存在、却无法在目标环境运行的代码。
05
output/driving_video_dataset_v1/├── driving_video_dataset_v1.py├── requirements.txt├── .env.example├── driving_video_dataset_v1_walkthrough.md└── driving_video_dataset_v1_submit.py
|
|
|
|
|
交付物不是一段需要用户继续补齐的示例代码,而是一套面向目标 Workspace 的作业包。
06
1. 以目标发布版本为准选择算子
EMR Daft 源码、开发环境和 Serverless 发布镜像可能存在版本差异。Skill 会优先核对当前 API;必要时,在获得授权后通过轻量探测任务确认目标镜像的实际导出,再生成主作业。
例如,目标镜像提供 ai_query 时,可以用它完成通用图片、视频或帧序列理解;旧镜像没有该 API 时,则切换到语义匹配的 VideoUnderstanding 等专用算子,而不是强行导入不存在的能力。
2. 大数据始终留在分布式执行路径
生成的作业通过 Ray 执行媒体处理和模型调用,并以 Paimon 表作为统一的数据 sink。原始视频、关键帧等媒体对象保留在 OSS,Paimon 表保存媒体 URI、采集元数据、中间结果、标签、Embedding 和运行状态。完整结果不会为了写出而被集中拉回 Driver,从而减少大规模路采数据处理中的单点内存风险。
3. 失败数据可定位、可重跑
对于支持结构化状态的 AI 能力,作业会保留:
content或业务结果;error与完成状态;model;token 用量;
输入数据的稳定标识。
单条媒体失败时,结果中尽可能保留对应错误,便于按行筛选与重跑。Skill 不会用“所有行都必须成功”掩盖真实数据中可能存在的坏 URI、损坏视频或无权访问对象。
4. 中间结果和最终结果都可验证
昂贵的预处理阶段可以物化为 Paimon 中间表,模型生成的标签、Embedding、处理状态和行级错误也分别写入 Paimon 表。每张表写入前都会检查目标 schema,避免历史表结构与当前输出不一致导致作业运行到最后才失败。
对于智驾长链路,验收会分别检查:
Ray 远程执行;
OSS 写入与回读;
模型结果及错误字段;
各类 Paimon 表的提交、记录数校验与原生回读;
最终组合 Pipeline 的 Driver 断言。
5. 企业安全约束默认开启
AccessKey、SecurityToken、模型凭据和签名 URL 不写入源码;
.env.example中所有 secret 保持为空;EMR 环境优先使用平台凭据链与挂载的 STS 凭据;
默认使用 EMR Model Manager 或平台批准的模型服务入口;
对外部模型可访问的媒体地址只在运行时签名,结果中保留原始
oss://标识;日志不打印完整提交响应、临时凭据或带签名的日志链接;
任务完成后自动释放任务级 Ray 集群。
07
静态检查和 --dry-run 只能证明交付物结构正确,不能替代目标环境验收。只有拿到 rj-* Submission ID、GetRayJob 到达 Succeeded、Driver 断言通过,并且目标数据回读成功后,Skill 才会报告远程验收通过。
如果用户尚未授权提交,或缺少 Workspace 权限和运行参数,结果会明确标记为 REMOTE_ACCEPTANCE_NOT_RUN,不会把“代码已经生成”包装成“作业已经跑通”。
08
智驾路采视频需要批量切分、抽帧、理解与结构化标注;
原始路采视频需要先完成转码、裁剪或压缩;
图片和视频需要生成 Embedding,用于相似检索、去重或样本挖掘;
昂贵的媒体预处理与模型阶段需要解耦和独立重跑;
中间结果、结构化标签、多模态向量、质量指标和错误记录需要统一写入 Paimon,形成可复用的数据资产;
已有 demo 需要补齐依赖、环境配置、提交脚本和远程验收;
团队希望把智驾数据 Pipeline 的工程约束固化为一致的生成规则。
它不会替代智驾团队定义场景 taxonomy、标注规范、Prompt、质量阈值和模型评测,但可以把这些业务定义稳定地转化成可运行、可提交、可验证的数据处理作业。
09
智驾数据处理真正消耗时间的,往往不是某一个视频算子或某一次模型调用,而是把媒体处理、分布式资源、模型服务、存储、安全、重跑和验收同时处理正确。
alibabacloud-emr-daft-multimodal-job 把这些重复的工程决策交给 Skill:
智驾团队定义数据、场景和质量目标;
Qoder 生成完整的 EMR Daft 作业包;
EMR Daft 在目标 Workspace 中分布式处理;
Skill 回读结果并给出可核验的验收结论。
现在可以在 Qoder 中输入一句:
“帮我生成一个智驾路采视频处理作业,完成场景切分、标签生成和向量化;原始媒体保留在 OSS,视频清单、中间结果、标签和向量统一写入 Paimon 表。”
从原始视频到可复用智驾数据资产,一次对话完成作业交付。
相关链接
EMR Serverless Spark 官方文档
https://help.aliyun.com/zh/emr/emr-serverless-spark/product-overview/what-is-emr-serverless-spark
Qoder CLI Skills 文档
https://docs.qoder.com/cli/Skills
Paimon 官方文档
https://paimon.apache.org/
Daft 官方文档
https://docs.daft.ai/en/stable/

/ END /
点击“阅读原文”了解更多关于大数据&AI产品解决方案~

