阿里云EMR Serverless Spark集成Ray:构建统一AI数据处理平台
1. 项目概述当Spark遇见Ray数据处理的新范式最近在折腾一个多模态数据分析的项目文本、图像、音频数据混在一起传统的Spark批处理框架在处理这种异构数据流时总感觉有些力不从心。批处理对图像特征提取这种计算密集型任务响应不够快而如果想引入一些实时的机器学习推理架构又会变得异常复杂。就在这个当口我注意到了阿里云EMR Serverless Spark的一个新动向——对Ray的深度集成与“再进化”。这可不是简单地把两个开源框架塞到一起而是真正从底层架构上为构建面向AI时代的数据处理“新基建”提供了一种全托管的、高性能的解决方案。简单来说阿里云EMR Serverless Spark全托管Ray核心解决的是在一个统一的、Serverless的数据平台上如何高效、灵活地处理批处理、流处理、机器学习训练、模型服务乃至更复杂的AI工作负载。Spark以其强大的批处理和SQL能力著称而Ray则专为分布式Python应用和AI负载设计原生支持Actor模型特别擅长处理有状态、低延迟的任务。将两者结合意味着你可以用Spark处理海量的结构化数据ETL然后无缝地将中间结果比如特征向量交给Ray集群进行实时的模型推理或强化学习训练整个过程无需在不同的集群或服务间搬运数据也无需操心底层资源的运维。这非常适合那些数据架构正在向AI化、实时化演进的中大型团队。比如你可以构建一个智能内容审核管道用Spark Streaming处理日志流、识别文本关键词同时将关联的图片和视频帧数据实时分发给Ray集群由部署在Ray上的深度学习模型进行并发的内容识别最后将结果统一写回数据湖。所有计算都在一个托管环境中自动伸缩你只需要关注业务逻辑。2. 核心架构解析Spark与Ray的深度融合之道2.1 为何是Spark Ray而非单一框架在深入技术细节前我们先理清一个根本问题为什么需要把Spark和Ray放在一起市面上优秀的计算框架很多比如Flink擅长流处理Dask在Python生态也有一定地位。选择SparkRay组合背后是计算范式互补的考量。Spark的核心抽象是RDD弹性分布式数据集其计算模型是面向数据集的、无状态的、函数式的转换。这种模型对于ETL、聚合分析等任务极其高效和稳定。然而当场景切换到需要长时间运行、有状态、且需要大量细粒度任务间通信的AI应用时例如一个模拟环境中的多个智能体需要持续交互Spark的模型就显得有些笨重。Ray的诞生正是为了弥补这一缺口。它的核心抽象是Actor和Task。Actor是一个有状态的工作进程可以常驻内存随时响应请求非常适合部署模型服务、模拟器或任何需要维护状态的服务。Task则是无状态的函数计算。Ray的动态任务图调度能力使得它能够极其灵活地处理大量细粒度的、依赖关系复杂的Python任务这正是现代AI工作负载的典型特征。在阿里云EMR Serverless Spark的集成方案中Spark和Ray并不是孤立运行的。其架构可以理解为“Spark作为数据面Ray作为计算面”的协同。Spark Job负责将数据从OSS、Hudi等数据源加载、清洗、转换成适合AI模型消费的格式例如将图片预处理成Tensor。然后它可以通过内置的集成通道将数据或任务直接提交到同一个虚拟集群VPC内共存的Ray集群中。Ray集群接收到这些任务后利用其分布式调度能力将任务分发到各个Actor或Worker上执行例如进行批量预测、模型微调或强化学习迭代。2.2 全托管Ray的“再进化”体现在何处“全托管”意味着阿里云帮你管理了Ray集群的生命周期自动创建、配置优化、监控告警、弹性伸缩以及故障恢复。而“再进化”则是在此基础上做了更深层次的优化和集成主要体现在以下几点资源隔离与共享的精细化早期的混部方案可能面临资源争抢问题。进化后的版本通过更精细的容器化隔离和资源配额管理确保Spark作业和Ray任务在共享底层资源池时互不干扰。Spark的Executor和Ray的Worker可以更安全地部署在同一组物理节点上提高整体资源利用率。数据交换零拷贝或高效序列化SparkJVM和RayPython之间最大的性能瓶颈往往是数据序列化与反序列化。阿里云对此进行了深度优化可能采用了Arrow内存格式作为中间层。Spark将数据转换为Arrow格式后Ray的Python进程可以直接通过共享内存或零拷贝方式读取避免了昂贵的JNI调用和Python Pickle序列化开销。这对于传输大型特征矩阵或嵌入向量至关重要。统一的监控与诊断界面将Spark UI和Ray Dashboard的关键指标进行了整合或关联。你可以在一个控制台上看到Spark作业的资源消耗、Stage进度同时也能观察到Ray集群的Actor状态、任务吞吐量、GPU利用率等提供了端到端的可观测性。与云原生服务的开箱即用集成除了Spark和Ray这套方案通常与阿里云的对象存储OSS、日志服务SLS、监控服务ARMS等深度集成。例如Ray任务的日志可以直接投递到SLS指标上报到ARMS使得运维体验与云上其他服务完全一致。注意这种深度集成并非没有代价。它在一定程度上将你绑定在了阿里云的技术栈上。虽然Spark和Ray本身是开源的但其中关于性能优化、网络互通、监控集成的“胶水代码”是阿里云的增值部分。评估时需权衡便利性与厂商锁定的风险。3. 实战构建一个多模态特征提取与匹配流水线理论说得再多不如动手试一下。我们以一个具体的场景为例假设我们有一个商品库包含商品标题文本、主图图像和宣传视频的音频轨音频。我们需要构建一个流水线能同时提取这三种模态的特征并实现“以图搜图”或“文本搜商品”的跨模态检索能力。3.1 环境准备与项目初始化首先你需要在阿里云EMR控制台创建一个Serverless Spark作业。与普通Spark作业不同我们需要在高级配置中启用Ray集成。创建作业在EMR Serverless控制台点击“创建作业”选择Spark引擎。依赖管理在“依赖”部分除了提交你的主应用JAR包或Python文件关键是要指定Ray相关的依赖。对于Python作业你可以在spark.submit.pyFiles中打包一个requirements.txt或者直接上传包含所有依赖的虚拟环境ZIP包。一个典型的requirements.txt可能包含ray[default]2.10.0 torch2.3.0 torchvision0.18.0 transformers4.40.0 librosa0.10.1 pillow10.3.0 pandas2.2.2 pyarrow15.0.2对于Java/Scala作业你需要确保Ray Java API的JAR包在classpath中。资源配置这是核心步骤。你需要在SparkConf中配置Ray集群的参数。通常这通过作业的“配置”项以spark.executorEnv.*或spark.ray.*为前缀的配置项来完成。# 启用Ray并设置启动参数 spark.executorEnv.RAY_ENABLE_WORKERtrue spark.ray.head.port6379 spark.ray.object-store-memory2000000000 # 2GB # 指定每个Executor内部启动Ray worker的资源 spark.executor.cores4 spark.executor.memory8g spark.executorEnv.RAY_WORKER_CPU2 spark.executorEnv.RAY_WORKER_MEMORY4096 # 4GB这段配置的意思是每个Spark Executor拥有4核8G在启动时会同时启动一个Ray Worker这个Worker被分配了2核4G的资源。剩下的资源2核4G留给Spark的JVM执行计算。这种“共置”模式减少了网络开销。3.2 核心代码拆解从Spark到Ray的接力我们的主程序假设是PySpark结构如下from pyspark.sql import SparkSession import ray import os def init_ray_on_spark(): 在Spark Executor中初始化Ray运行时 # 从环境变量获取Ray head节点的地址这个地址由EMR在启动Ray集群时自动设置 ray_head_ip os.environ.get(RAY_HEAD_SERVICE_HOST, localhost) ray_head_port os.environ.get(RAY_HEAD_SERVICE_PORT, 6379) ray_address f{ray_head_ip}:{ray_head_port} # 连接到已有的Ray集群由EMR管理 if not ray.is_initialized(): ray.init(addressray_address, ignore_reinit_errorTrue, runtime_env{working_dir: ., pip: [torch, transformers, librosa]}) print(fRay initialized on worker, address: {ray_address}) # 1. 初始化SparkSession并确保Ray相关配置生效 spark SparkSession.builder \ .appName(MultimodalFeaturePipeline) \ .config(spark.sql.execution.arrow.pyspark.enabled, true) \ # 启用Arrow加速 .getOrCreate() # 2. 注册一个UDF这个UDF会在每个Executor上执行用于初始化Ray spark.sparkContext.parallelize([1]).foreachPartition(lambda _: init_ray_on_spark()) # 3. 使用Spark读取原始数据例如一个包含图片OSS路径、文本、音频路径的Parquet表 df spark.read.parquet(oss://your-bucket/data/product.parquet) # 4. 定义Ray远程函数Actor用于特征提取 ray.remote(num_gpus0.5) # 假设每个任务需要0.5个GPU如果纯CPU则去掉此参数 class ImageFeatureExtractor: def __init__(self, model_nameclip-vit-base-patch32): from transformers import CLIPProcessor, CLIPModel self.device cuda if ray.get_gpu_ids() else cpu self.model CLIPModel.from_pretrained(model_name).to(self.device) self.processor CLIPProcessor.from_pretrained(model_name) def extract(self, image_url): # 从OSS下载图片或直接读取字节流假设已配置好OSS访问权限 image_data download_from_oss(image_url) inputs self.processor(imagesimage_data, return_tensorspt).to(self.device) with torch.no_grad(): image_features self.model.get_image_features(**inputs) return image_features.cpu().numpy().flatten().tolist() # 返回列表便于Spark存储 ray.remote class TextFeatureExtractor: def __init__(self, model_namesentence-transformers/all-MiniLM-L6-v2): from sentence_transformers import SentenceTransformer self.model SentenceTransformer(model_name) def extract(self, text): return self.model.encode(text).tolist() # 5. 在Spark DataFrame操作中调用Ray Actor # 首先获取Ray Actor的句柄。注意Actor创建本身是远程操作。 # 一种模式是广播Actor句柄但更常见的做法是在mapPartitions内部分别创建以利用数据本地性。 def extract_features_partition(iterator): # 在每个分区内初始化Ray连接和Actor init_ray_on_spark() image_extractor ImageFeatureExtractor.remote() text_extractor TextFeatureExtractor.remote() results [] for row in iterator: # 同步调用Ray任务并等待结果。对于生产环境可考虑异步批量提交以提高吞吐。 img_feat_future image_extractor.extract.remote(row[image_url]) txt_feat_future text_extractor.extract.remote(row[title]) img_feat ray.get(img_feat_future) txt_feat ray.get(txt_feat_future) # 将提取的特征作为新列返回 results.append((row[product_id], row[title], row[image_url], img_feat, txt_feat)) return results # 由于Ray操作不能在标准的Spark SQL UDF中直接进行我们使用RDD的mapPartitions result_rdd df.rdd.mapPartitions(extract_features_partition) # 6. 将结果转换回DataFrame并保存 schema ... # 定义包含特征向量的Schema result_df spark.createDataFrame(result_rdd, schemaschema) result_df.write.mode(overwrite).parquet(oss://your-bucket/features/product_features.parquet) # 7. 关闭资源在Serverless环境下通常会自动清理 spark.stop()3.3 性能调优与资源配置心得在实际部署中直接使用上述代码可能会遇到性能问题。以下是几个关键的调优点Actor复用 vs. 任务粒度在上面的例子中我们为每个分区创建了一组Actor。如果分区内数据量很大例如数万条这是高效的因为Actor被重用了。但如果分区很小频繁创建销毁Actor的开销就很大。你需要根据数据量调整Spark的spark.sql.shuffle.partitions或repartition来控制分区大小确保每个分区有足够的工作量让Ray Actor“吃饱”。批处理Batching同步调用ray.get一条一条处理是低效的。更好的模式是让Actor的extract方法接受一个列表进行批量处理。例如ImageFeatureExtractor.extract.remote([url1, url2, ...])。这能极大减少远程调用的开销和GPU的启动延迟。资源死锁预防Ray集群的资源是全局调度的。如果你提交了多个需要GPU的Actor任务但总GPU需求超过了物理GPU数量任务会排队等待可能导致整个流水线卡住。务必在ray.remote中准确声明资源需求如num_gpus0.5并确保你的EMR作业申请了足够的GPU总量。同时可以考虑使用Ray的placement_group进行资源预留确保关键任务能获得资源。数据本地性尽可能让Ray Worker处理存储在同一台机器或同一可用区上的数据。在EMR环境中如果数据在OSS可以通过OSS的VPC内网端点访问来降低延迟。对于更大的数据集可以考虑先用Spark将数据预处理并缓存到集群本地磁盘如ESSD再由同节点的Ray Worker读取。4. 深入原理Spark与Ray间的数据高效流转为什么这种架构能加速关键在于数据交换层。传统做法可能是Spark将结果写入HDFS或OSS然后另一个Ray作业再去读取这涉及磁盘I/O和网络传输。在EMR Serverless Spark的集成环境中优化主要在两个层面Arrow内存格式作为通用语言Spark通过PySpark可以将DataFrame或RDD的数据以Apache Arrow的格式直接保存在堆外内存中。Arrow是一种列式内存格式被PythonPyArrow、Java、C等语言广泛支持。Ray的WorkerPython进程可以通过pyarrow库以近乎零拷贝的方式从同一块内存区域读取这些数据避免了JVM到Python的序列化成本。在代码中启用spark.sql.execution.arrow.pyspark.enabled配置就是为此。共享文件系统或对象存储的优化访问当数据量巨大无法全部放在内存中交换时中间数据会写入一个共享存储如OSS。EMR对此路径的访问做了优化可能是通过FUSE挂载或专用的SDK使得Spark和Ray都能以高性能并行读写。更重要的是由于它们在同一个VPC和安全组下网络延迟和带宽都得到保障。下表对比了不同数据交换方式的性能特点交换方式原理优点缺点适用场景Arrow内存零拷贝Spark将数据转为Arrow格式存于堆外内存Ray直接访问。延迟极低吞吐量高无序列化开销。受限于Executor内存大小不适合超大单批数据。中小型特征数据集、迭代计算中的中间结果。Parquet/OSS交换Spark写结果为Parquet到OSSRay从OSS读取。容量无限数据持久化兼容性好。有磁盘I/O和网络延迟速度较慢。最终结果存储、跨作业/天级别的数据交换。Ray内部对象存储Spark通过Ray的API将数据存入Ray的分布式对象存储。Ray任务访问快支持复杂对象。需要额外的集成代码对象存储容量有限。Spark预处理后供多个Ray任务反复使用的共享数据。在实际操作中我通常会采用混合策略对于流水线中的高频、小批量特征传递尽量设计成通过Arrow内存交换对于最终的产出或检查点则规整地写入OSS Parquet。5. 常见问题与故障排查实录即便有了全托管服务在实际开发中依然会遇到各种问题。下面记录几个我踩过的坑和解决方法。5.1 Ray Worker启动失败或无法连接Head节点这是最常见的问题。通常体现在Spark作业日志中大量出现Connection refused或Timeout错误。排查思路1检查网络与安全组。确保你的EMR Serverless Spark作业和其创建的Ray集群在同一个VPC内。Ray的Head节点默认监听端口6379需要确保Spark Executor所在的安全组允许出站连接到Head节点的IP和此端口。在全托管环境下阿里云通常会自动配置好但如果你自定义了网络这里是重点检查对象。排查思路2检查资源冲突。spark.executorEnv.RAY_WORKER_CPU和RAY_WORKER_MEMORY的设置不能超过spark.executor.cores和spark.executor.memory。例如Executor总内存8g你给JVM堆内存设置了6g通过spark.executor.memoryOverhead和spark.executor.memory调节那么留给Ray Worker的内存包括堆外就必须小于2g否则启动时会因内存不足失败。排查思路3查看Ray Head日志。在EMR控制台找到对应的Ray集群如果有独立入口或Spark Driver日志搜索ray.init相关的错误。常见的错误是Python环境不兼容比如Ray版本与pickle协议不匹配。确保所有节点上的ray版本完全一致。5.2 任务执行慢或GPU利用率低现象是Spark作业整体运行时间远超预期通过监控发现Ray集群的GPU利用率波动很大或一直很低。排查思路1任务粒度太细。如果你有10万条数据创建了10万个Ray远程任务每个任务只处理一张图片那么大量的时间都花在了任务调度和启动上而不是计算上。解决方案使用mapPartitions并在分区内部进行批处理或者使用Ray的ray.wait进行异步批量提交将一个批次的数据比如100条作为一个任务提交。排查思路2数据加载成为瓶颈。特征提取代码中download_from_oss(image_url)如果是同步且单线程的那么大部分时间可能在等待网络I/O。解决方案使用异步HTTP客户端如aiohttp在Actor内部并发下载或者更优的是在Spark侧先将图片批量预取到本地SSD然后Ray Worker从本地读取。排查思路3Actor初始化开销大。每次处理都新建一个加载了大型模型的Actor模型加载可能耗时几十秒。解决方案使用Ray的Actor Pool模式或者确保Actor是长生命周期的在一个分区内被反复调用。在ray.remote中设置max_concurrency参数可以让一个Actor同时处理多个请求提高吞吐。5.3 内存溢出OOM问题OOM可能发生在Spark侧JVM OOM也可能发生在Ray侧Python进程被Kill。Spark JVM OOM通常是因为Arrow转换的数据量过大或者Spark SQL操作如collect拉取了过多数据到Driver。解决增加Executor内存spark.executor.memory并适当增加堆外内存spark.executor.memoryOverhead。更重要的是避免在Driver端收集大量数据尽量用分布式操作。Ray Worker OOMPython进程因处理大量数据或模型本身占用显存/内存过多而被系统终止。解决在ray.remote中明确设置num_cpus和memory参数让Ray调度器更合理地分配资源。例如ray.remote(num_cpus2, memory4*1024*1024*1024)表示请求2核和4GB内存。检查模型加载是否占用了过多显存。对于多GPU任务使用CUDA_VISIBLE_DEVICES环境变量或Ray的num_gpus参数来隔离。在代码中及时释放不再需要的Tensor或大对象可以调用del并手动触发垃圾回收gc.collect()。5.4 如何调试和监控混合任务全托管的好处是提供了集成的监控视图但你仍需知道看哪里。Spark UI关注在“Stages”页找到你调用mapPartitions的那个Stage。看其执行时间、GC时间、Shuffle数据量。如果这个Stage特别慢问题很可能出在Ray任务上。Ray Dashboard这是最重要的工具。EMR通常会提供一个代理地址来访问Ray Dashboard。在这里你可以看到Cluster总CPU/GPU/内存使用情况节点状态。Jobs提交的所有Ray任务它们的状态Pending、Running、Finished、执行时间。Actors所有Actor的状态、所在节点、资源使用情况。如果Actor不断重启说明可能遇到了崩溃。Logs查看具体Actor或任务的stdout/stderr日志这是定位Python错误的关键。阿里云日志服务SLS如果配置了日志收集你可以在SLS中按作业ID或Ray Job ID搜索聚合查看所有节点的日志便于全局排查。6. 演进思考从特征工程到在线服务的闭环这套方案的价值不止于离线的特征提取。它的真正威力在于能够平滑地支撑起从数据预处理、模型训练到在线服务的完整AI闭环。设想这样一个场景你利用SparkRay高效地完成了历史数据的特征提取和模型训练。训练好的模型例如一个PyTorch模型文件被保存到了OSS。接下来你可以几乎不做任何架构改动就将这个模型部署为在线服务。具体做法是在Ray集群中启动一个模型服务Actor。这个Actor在初始化时从OSS加载训练好的模型。ray.remote(num_gpus0.5) class ModelServingActor: def __init__(self, model_path): self.model load_model_from_oss(model_path) self.model.eval() def predict(self, input_features_batch): with torch.no_grad(): predictions self.model(input_features_batch) return predictions.cpu().numpy().tolist()然后你可以通过Ray的serve模块或者简单地暴露一个HTTP端点使用FastAPI集成让外部的应用比如一个推荐系统的API服务器直接调用这个Actor进行实时预测。由于Ray Actor是常驻的预测的延迟可以做到毫秒级。与此同时流式的Spark作业可以持续将实时产生的数据如用户点击流转化为特征并调用这个在线的Ray Actor进行实时打分实现流批一体的AI推理。所有计算资源仍然由同一个EMR Serverless Spark环境弹性调度运维复杂度并没有增加。这种模式打破了传统Lambda架构中批处理层和速度层流处理层的壁垒也模糊了数据平台和机器学习平台之间的界限。对于追求敏捷和效率的团队来说它提供了一个极具吸引力的统一计算底座。当然这要求团队对Spark和Ray都有一定的掌握深度并且充分理解云原生的资源管理方式。