大数跨境

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

阿里云 EMR Serverless Spark 全托管 Ray 再进化:加速构建全模态数据处理新基建 阿里云大数据AI平台
2026-07-20
4
导读:阿里云 EMR Serverless Spark + Ray 双引擎构建全模态数据处理的新基建,通过极致内核优化和统一数据、算力底座,彻底打通了大数据工程与 AI 模型训练的割裂。

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

01

Data + AI 深度融合,呼唤全链路计算底座

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

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

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

02

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

产品演讲路径

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

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

图片

03

Spark × Ray 双引擎架构

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

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

图片

除标准 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、Data-Juicer 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 稳定高效支撑了多种工作负载,包括视频切割、抽帧、图片打标、图片向量化、文本生成等场景。

04

统一底座:一份数据、一池算力、无缝融合

统一数据底座

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

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

统一算力与基础设施

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

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

图片

Spark × Ray 融合计算

Spark  Ray 的融合计算有两种方式:

  • 管线串联作为数据管线的不同节点,互为依赖处理不同模态的数据,以持久化的表或 Raw Data 作为传输介质。

  • RayDP 零拷贝依赖 RayDP  Spark 运行在 Ray 上,通过 Ray Object Store 实现 Spark Dataframe  Ray Dataset 之间的零拷贝互转,从而在一个作业里同时运行 Spark  Ray。

05

行业落地案例

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=2if 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)

案例二:具身智能 Egocentric 视频抽帧筛选

训练具身智能模型理解和复现人类日常操作动作时,研发人员通常会让真实人类佩戴头戴式摄像头(Head-mountedCamera),在厨房等真实生活场景中执行烹饪、切菜、翻炒等操作,采集大量第一人称视角(Egocentric)的原始视频。这类视频记录了人类在厨房中翻炒食材的完整过程,是训练VLA(视觉-语言-动作)模型的宝贵数据来源。

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

代码示例

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

import daftfrom daft import colfrom daft.functions import encode_imagefrom daft.emr.functions import emr_udffrom daft.emr.functions import ai_queryfrom 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. 编码为 JPEGdf = 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()

06

结语

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

参考文档:

  • 阿里云 EMR Serverless Spark 2026-04-15 功能发布记录:

https://help.aliyun.com/zh/emr/emr-serverless-spark/product-overview/2026-04-15-version

  • 创建 Ray 集群:

https://help.aliyun.com/zh/emr/emr-serverless-spark/developer-reference/api-emr-serverless-spark-2023-08-08-createraycluster

  • 向 Ray 集群提交任务:

https://help.aliyun.com/zh/emr/emr-serverless-spark/user-guide/submit-a-task-to-the-ray-cluster

图片

【声明】内容源于网络
0
0
阿里云大数据AI平台
阿里云大数据AI平台依托阿里领先的云基础设施、大数据和AI工程能力、场景算法技术和多年行业实践,一站式地为企业和开发者提供云原生的大数据和AI能力体系。帮助提升AI应用开发效率,促进AI在产业中规模化落地,激发业务价值。
内容 741
粉丝 0
阿里云大数据AI平台 阿里云大数据AI平台依托阿里领先的云基础设施、大数据和AI工程能力、场景算法技术和多年行业实践,一站式地为企业和开发者提供云原生的大数据和AI能力体系。帮助提升AI应用开发效率,促进AI在产业中规模化落地,激发业务价值。
总阅读7.6k
粉丝0
内容741