阿里云 EMR Serverless Spark 全托管 Ray 再进化:加速构建全模态数据处理新基建

简介: 阿里云 EMR Serverless Spark + Ray 双引擎构建全模态数据处理的新基建,通过极致内核优化和统一数据、算力底座,彻底打通了大数据工程与 AI 模型训练的割裂。结合 RayData、Daft、Data-Juicer 等多模态引擎,以及 CPFS、OSS 等高性能存储生态,阿里云正在为全球的 AI 开发者提供一套最具竞争力的数据新基建。

阿里云 EMR Serverless Spark 自今年4月推出全托管 Ray 以来,能力持续升级,在同一工作空间内深度集成 Spark 与 Ray,以统一湖仓存储、统一 CPU&GPU 资源池、统一安全和运维体系,承载从结构化数据处理到 AI 训练数据预处理、后训练、推理与服务的完整链路。目前,该能力已在具身智能、智能制造、自动驾驶等行业实现生产级落地,支撑企业在一套云原生平台上高效完成数据治理、样本加工、模型开发和规模化生产。


数据与 AI 深度融合,需要贯通全链路的计算底座

生成式 AI 正在重塑企业的数据处理范式。传统数据工程主要面向结构化表、日志和文本;而具身智能、自动驾驶等新一代 AI 系统工程,则持续产生第一视角视频、连续图像、语音、传感器记录和控制轨迹。除了数据规模扩大,数据形态、算子类型和资源需求也在发生深刻变化:一条生产流水线中,既包含大规模 SQL、Join、去重和统计,也涵盖视频解码、图像理解、向量化、样本过滤、模型推理等 Python 原生任务,并同时消耗 CPU、GPU、网络和存储带宽。


具身智能中的 Egocentric 数据是典型代表。机器人或可穿戴设备以第一视角记录人与环境的持续交互,单个样本往往同时包含视频片段、语音、文本指令、动作序列和时间戳。自动驾驶同样需要围绕海量路测数据完成切帧、场景识别、长尾目标发现、标注、去重、特征生成和回灌。当前大量视觉、语音、强化学习和大模型工具均围绕 Python 生态构建,这些任务需要可横向扩展的 Python 分布式运行环境;与此同时,样本目录管理、批量 ETL、数据质量、湖表写入和跨任务治理,仍然依赖成熟的大数据处理能力。


为此,阿里云 EMR Serverless Spark 全托管 Ray 不断优化协同能力:Spark 负责规模化数据处理、SQL 分析与湖仓读写,Ray 负责 Python 原生分布式计算、多模态数据处理、训练推理与模型服务。二者面向同一业务流程无缝协同,共同构成 Data + AI 一体化计算底座,为用户提供面向数据与 AI 工作负载的全托管 Lakehouse 引擎。


架构演进:极致全模态内核优化,统一数据和算力底座

产品演讲路径

从正式上线开始,产品围绕可靠性、Ray内核优化、Spark 与 Ray 协同,以及统一数据和算力底座持续增强。6 月进一步增加 Ray Job、History Server 等能力,完善批处理和常驻集群两种使用形态。


今天,Spark 与 Ray 已在同一工作空间中共享湖仓数据、CPU/GPU 资源池、安全和运维体系,共同支撑从数据治理、样本加工到训练与推理的完整链路。


时间

产品特性

2026-04

EMR Serverless Spark 正式发布 Ray 集群资源形态,支持创建、管理 Ray 集群及 GPU、弹性、安全能力。

2026-05

持续增强 Ray 稳定性、Spark/Ray 协同以及统一数据和算力底座。

2026-06

增加 Ray Job、History Server 等产品能力,完善批处理和常驻集群两种提交形态。

当前阶段

面向全模态数据处理、训练数据准备与模型服务,持续完善 Data + AI 一体化架构。


EMR Spark:高性能数据工程与湖仓处理

EMR Serverless Spark的结构化执行引擎Fusion2.0在TPC-DS 100T官方Benchmark取得世界第一的成绩,相比前榜首性能提升100%,性价比提升500%。

b378e7ac85714f42a267790d2369be6d.png

除标准 Benchmark 外,Fusion 2.0 基于海量生产实践,围绕 I/O、数据倾斜、半结构化数据解析、大规模 Shuffle、磁盘溢写及历史信息利用等场景开展 QO/QE 联合优化,使大量真实作业获得数倍性能提升。


EMR Ray: 稳定高性能的多模态处理与推理

Ray是全栈AI计算框架,包含Ray Core和Ray AI Libraries,涵盖资源管理、高性能存储、计算原语、多模态数据处理、AI训练、AI推理服务等重要功能。


EMR Ray Core做了诸多优化提升稳定性,包括Ray Head高可用,多级节点容错,基于磁盘水位的自动扩缩容,GCS外挂高可用Redis,坏卡自动检测隔离等。


EMR Ray的多模态处理能力主要来自Ray Data,Daft on Ray,Datajuicer on Ray,以及丰富的python多模态处理生态。多模态处理和结构化处理并非割裂,而是侧重点不同,结构化处理侧重关系算子,而多模态侧重图片、文本、音视频的原生类型支持、GPU加速推理、大模型推理等。


EMR Ray & Daft性能优化包含三方面。一是拓展Fusion边界,把关系算子的优化复制到Ray Data和Daft,包括向量化算子,对接Celeborn,Query Optimizer,Pipeline多线程,磁盘溢写等。二是提升GPU利用率,方法包括自动拆解和异步化CPU和GPU算子/UDF,自动扩容CPU避免GPU饥饿,自适应显存超卖等。三是提升调用大模型服务的优化,包括异步化,Batch化,多层次QO优化等,同时降低Token消耗和作业e2e延迟。


EMR Ray提供Ray Cluster和Ray Job两种作业提交形态,前者复用常驻集群最大化降低冷启动时间,适用于环境复用、低延迟场景;后者按作业申请资源,适用于环境不统一、稳定性要求高的批处理场景。


EMR Ray稳定高效支撑了多种工作负载,包括视频切割,抽帧,图片打标,图片向量化,文本生成等场景。


Spark & Ray 统一数据底座

Spark和Ray构建在统一的数据管理底座之上,既支持开源开放的Hive Metastore(HMS) + OSS架构,也支持全托管DLF方案,既覆盖Paimon,Iceberg,Delta,Hudi等分钟级新鲜度的湖格式,也支持Fluss这种新兴的秒级流存储。


数据底座为Spark和Ray提供统一的全模态数据存储和管理服务,避免数据冗余。以Paimon为例,结构化和多模态数据分别以结构化和Blob类型存储,向量数据以Vector类型存储,同时提供标量和向量索引文件。Spark和Ray在同一份Paimon数据上读取、加工、生成、检索数据,互为上下游,真正做到一份数据,多引擎平权。


Spark & Ray 统一算力和基础设施

Spark和Ray共用一份算力池,工作空间的CPU和GPU算力按需、细粒度在两个引擎之间统一高效供应,Spark释放的资源能立即被Ray消费,反之亦然,从而消除资源碎片,最大化利用率。


除了算力池,Spark和Ray还共享其他基础设施,包括统一的认证鉴权体系,统一的Cache服务,挂载CPFS、NAS、OSS的能力,监控报警等。整体架构如下所示:

ca8b681b7662491bbe2142b1a5587b54.png


Spark & Ray 融合计算

Spark和Ray的融合计算有两种方式。一是作为数据管线的不同节点,互为依赖处理不同模态的数据,并以持久化的表或Raw Data作为传输介质。二是依赖RayDP,让Spark运行在Ray上,通过Ray Object Store实现Spark Dataframe和Ray Dataset之间的零拷贝互转,从而在一个作业里同时运行Spark和Ray。


案例

Ray发布以来,在具身智能、智能制造、自动驾驶等行业落地多个客户。


多模态数据处理并非单一算子问题。例如,一个视频数据集往往需要经历解析、切分、质量过滤、内容去重、语义标注、向量化和格式转换。全托管 Ray 为 RayData、Daft、Data-Juicer 等面向 AI 数据而生的引擎与工具提供统一运行底座,使企业能够在不自建 Ray 基础设施的前提下使用 Python 多模态生态。


案例一

自动驾驶路测会产生大量摄像头图片,其中绝大多数是正常道路场景。该流水线直接读取 OSS 中的路测图片,过滤低质量画面,再调用千问多模态模型识别施工区域、行人横穿、事故车辆、救护车及恶劣天气等长尾场景,最终将标注结果直接写回 OSS。


代码示例

Ray Data 是面向 AI 与多模态数据处理的分布式数据引擎,可通过 Python API 将图片读取、质量过滤和大模型推理等 I/O、CPU 与 GPU 算子组织为流水线,并由 Ray 统一完成资源调度、并发控制与背压管理。EMR Ray 原生支持读写 OSS,可从 OSS 并行读取海量图片,完成处理后将标注结果分布式写回 OSS,简化多模态数据流水线的开发与运维。

INPUT_PATH = "oss://<bucket>/autonomous/camera_frames/"
OUTPUT_PATH = "oss://<bucket>/autonomous/scene_labels/"

PROMPT = """
判断这张自动驾驶路测图片属于哪个场景。
只返回以下一个标签,不要输出解释:
NORMAL、PEDESTRIAN_CROSSING、CONSTRUCTION、
EMERGENCY_VEHICLE、ACCIDENT、BAD_WEATHER。
"""

ray.init()

# 使用 CPU 并行过滤分辨率过低、过暗或低对比度图片。
def filter_image_quality(batch: pd.DataFrame) -> pd.DataFrame:
    keep = []

    for image in batch["image"]:
        height, width = image.shape[:2]
        gray = image.mean(axis=2) if image.ndim == 3 else image

        brightness = float(gray.mean())
        contrast = float(gray.var())

        keep.append(
            width >= 1280
            and height >= 720
            and 25 <= brightness <= 230
            and contrast >= 100
        )

    return batch.loc[keep].reset_index(drop=True)


start = time.time()

# Ray Data 直接读取 OSS 图片。
ds = ray.data.read_images(
    INPUT_PATH,
    include_paths=True,
    file_extensions=["jpg", "jpeg", "png"],
)

# CPU 图片质量过滤。
ds = ds.map_batches(
    filter_image_quality,
    batch_format="pandas",
    batch_size=64,
    concurrency=32,
    num_cpus=1,
)

# EMR Ray 内置大模型算子:
# 将 OSS 图片路径直接提交给百炼多模态模型,并自动管理批次与并发。
ds = ai_query(
    ds,
    prompt=PROMPT,
    data_column="path",
    output_column="scene_label",
    model="qwen3.7-plus",
    concurrency=16,
    batch_size=8,
    options={"enable_thinking": False},
)

# 移除解码后的图片数据,只保存路径和模型标签。
ds = ds.select_columns(["path", "scene_label"])

# Ray Data 分布式直接写回 OSS。
ds.write_parquet(OUTPUT_PATH)

print(ds.stats())
print("Runtime:", time.time() - start)


案例二

训练具身智能模型理解和复现人类日常操作动作时,研发人员通常会让真实人类佩戴头戴式摄像头(Head-mounted Camera),在厨房等真实生活场景中执行烹饪、切菜、翻炒等操作,采集大量第一人称视角(Egocentric)的原始视频。


这类视频记录了人类在厨房中翻炒食材的完整过程,是训练 VLA(视觉-语言-动作)模型的宝贵数据来源。


这些视频中包含大量无价值的冗余帧——例如操作间隙的静止画面、模糊镜头,或镜头朝向地面/天花板时拍摄到的无效内容。我们需要一个高效的自动化流水线,从海量视频中抽取关键帧,并利用多模态大模型对图片进行智能筛选,最终只保留包含有效烹饪操作内容(例如"锅具"或"食材处理")的高质量图片帧,用于后续的模型训练。


代码示例

Daft 提供 Python DataFrame 与 SQL 接口,可将视频解码、图像处理、模型推理等 CPU/GPU 算子组织为流水线,并由 Ray 负责分布式资源调度。其面向多模态数据的流式执行与背压机制,有助于让 I/O、CPU 预处理和 GPU 推理并行衔接,并支持 Paimon 等开放数据格式。

import daft
from daft import col
from daft.functions import encode_image
from daft.emr.functions import emr_udf
from daft.emr.functions import ai_query
from daft import lit

# 流式读取 OSS 视频,仅提取关键帧,每2秒采样一次
df = daft.read_video_frames(
    path="oss://<your-bucket>/cooking_videos/*.mp4",
    image_height=480,
    image_width=640,
    is_key_frame=True,           # 只取关键帧,减少冗余
    sample_interval_seconds=2.0, # 采样频率
    max_frames_per_video=100,    # 限制单视频最大帧数
    io_config=io_config,
)

# 1. 编码为 JPEG
df = df.with_column("jpeg_bytes", encode_image(col("data"), "JPEG"))

# 2. 批量上传至 OSS (并发控制最大化 I/O 吞吐)
df = df.with_column(
    "oss_path",
    emr_udf(
        FrameUploader,
        construct_args={"output_base": output_base},
        concurrency=4,
        batch_size=32,
    )(col("path"), col("frame_index"), col("jpeg_bytes")),
)

# 调用 Qwen 模型进行视觉判断
df = df.with_column(
    "llm_result",
    ai_query(
        prompt=lit("判断这张图片是否包含有效的烹饪操作画面(如锅具、食材处理、手部操作等)。如果包含请只回复KEEP,不包含请只回复DROP,不要输出其他内容"),
        data=col("oss_path"),      # 传入 OSS 路径,模型服务端直接拉取
        model="qwen3.6-plus",
        concurrency=32,            # 32并发,充分利用模型吞吐
        batch_size=32,
        options={"enable_thinking": False},
    ),
)

from daft.functions import get as struct_get

df = df.with_column("llm_content", struct_get(col("llm_result"), "content"))
df = df.with_column("keep", col("llm_content").upper().contains("KEEP"))

# 物化结果,统计保留与删除数量
df = df.collect()
result_dict = df.to_pydict()
total = len(result_dict["keep"])
kept = sum(1 for k in result_dict["keep"] if k)
dropped = total - kept

# 删除 LLM 判定为 DROP 的帧
if dropped > 0:
    df_drop = df.where(col("keep") != True).select("oss_path")
    df_drop = df_drop.with_column(
        "deleted",
        emr_udf(OSSDeleter, concurrency=4, batch_size=32)(col("oss_path")),
    )
    df_drop.collect()


结语

阿里云 EMR Serverless Spark + Ray 双引擎构建全模态数据处理的新基建,通过极致内核优化和统一数据、算力底座,彻底打通了大数据工程与 AI 模型训练的割裂。结合 RayData、Daft、Data-Juicer 等多模态引擎,以及 CPFS、OSS 等高性能存储生态,阿里云正在为全球的 AI 开发者提供一套最具竞争力的数据新基建。



参考文档:

目录
相关文章
|
29天前
|
人工智能 分布式计算 Serverless
EMR Serverless Daft 如何简化多模态数据处理:视频抽帧、清洗、标注全流程与具身智能实践
阿里云 EMR Serverless Spark 引入 Ray 分布式计算框架与 Daft 高性能数据引擎,为用户提供了一套开箱即用、免运维且极致高效的多模态数据处理基础设施。
|
23天前
|
存储 人工智能 Apache
从向量存储到 Agentic 数据基础设施:Paimon × Milvus 如何构建 AI 原生多模态数据湖
本文整理自李钰在 Apache Flink Forward Asia 2026 的演讲,讨论 AI 与 Agent 进入生产环境后,数据湖与向量数据库“双系统”架构暴露出的结构性问题,以及 Apache Paimon 与 Milvus 围绕同一份湖数据协同的技术路径。
从向量存储到 Agentic 数据基础设施:Paimon × Milvus 如何构建 AI 原生多模态数据湖
|
6月前
|
存储 分布式计算 数据建模
淘宝闪购基于阿里云 EMR Serverless Spark&Paimon的湖仓实践:超大规模下的特征生产&多维分析双提效
本文介绍阿里云 Serverless Spark + Paimon 在淘宝闪购大数据湖仓场景的应用。
|
23天前
|
SQL 人工智能 缓存
阿里云 EMR Serverless StarRocks(Stella 2.2.0)发布:多模态处理与分析闭环,内表与湖表统一检索
Stella 2.2 面向 AI 时代的数据基础设施,打通“多模态数据处理—向量化与理解—多路检索—分析消费”的完整闭环。无论数据沉淀在 Paimon 湖表,还是 StarRocks 存算分离内表,都可以在统一 SQL 入口下组合结构化分析、全文检索、向量检索与 AI Function,服务智能驾驶、具身智能、内容与商品理解、企业知识库和 RAG 等场景。
314 2
|
30天前
|
SQL 人工智能 Serverless
从数据湖到多模态湖仓-基于阿里云 EMR Serverless StarRocks 与 DLF Paimon 构建AI时代的统一分析检索架构
阿里云 EMR Serverless StarRocks 在统一数据、一致语义和系统级优化之上,构建了面向 AI Data、AI Agent 和多模态应用的下一代湖仓架构。
从数据湖到多模态湖仓-基于阿里云 EMR Serverless StarRocks 与 DLF Paimon 构建AI时代的统一分析检索架构
|
30天前
|
人工智能 分布式计算 DataWorks
阿里云大数据 AI 产品月刊-2026年6月
阿里云大数据& AI 产品技术月刊【2026 年 6 月】,涵盖 6 月技术速递、产品和功能发布、市场和客户应用实践等内容,帮助您快速了解阿里云大数据& AI 方面最新动态。
|
7月前
|
分布式计算 Serverless 测试技术
有奖实践:EMR Serverless StarRocks × Serverless Spark x DLF 共探 TPC 极致性能
免费试用 EMR Serverless StarRocks 与 EMR Serverless Spark,体验“实时分析冠军”与“批处理之神”的极致性能表现!
有奖实践:EMR Serverless StarRocks × Serverless Spark x DLF 共探 TPC 极致性能
|
29天前
|
自然语言处理 运维 算法
基于阿里云 OpenSearch 行业算法版的海外电商智能搜索实践
本文档详述东南亚电商巨头GoTerra从自建Elasticsearch迁移至阿里云OpenSearch行业算法版的全链路实践,涵盖多语言语义理解、向量检索、低延迟优化(P99延迟降幅超70%)、弹性成本控制与安全加固,提供可复用的企业级智能搜索升级范式。
222 0
|
9月前
|
存储 缓存 Cloud Native
EMR StarRocks Stella 内核正式发布,登顶 TPC 榜单全球第一
EMR Serverless StarRocks 重磅发布全新企业级版本内核 Stella (StarRocks Efficient and Lightening-fast Lakehouse),完全兼容开源 StarRocks,为用户提供企业级的产品功能、卓越的性能及稳定性保障。

热门文章

最新文章