ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

EMR Serverless Spark实战:CPU+GPU异构计算配置与调优

EMR Serverless Spark实战:CPU+GPU异构计算配置与调优 先说结论EMR Serverless Spark跑CPUGPU异构计算不是什么高不可攀的黑科技本质就是把“CPU算数据、GPU算模型”这两件事塞进同一个弹性资源池里让Spark的调度器自己去分活。你不需要单独搭K8s集群也不需要自己维护GPU驱动只需要在Serverless环境里把资源配置和作业参数写对就能让ETL和AI推理在同一套流程里跑起来。这篇文章我打算从“为什么要用异构”“怎么配才能跑起来”“调优和计费怎么平衡”“踩坑实录”四个角度展开尽量把配置参数和底层逻辑都讲透。无论你是数据平台工程师还是算法工程师想摆脱“手工传表给数仓”的尴尬都值得花几分钟看完。1. 内容整体设计与思路拆解1.1 为什么EMR Serverless Spark需要CPUGPU异构先说个很常见的场景业务那边每天要处理上亿条用户行为日志先用Spark做清洗和特征工程然后把这些特征喂给一个深度学习模型做实时打分最后把打分结果再写回在线存储。传统做法是分成两套系统Spark集群负责ETLGPU集群负责推理中间靠接口或者消息队列对接。这个方案最大的问题不是技术难度而是资源利用率。ETL任务通常有早晚高峰GPU任务也有自己的波峰波谷两套集群各自按峰值申请资源绝大多数时间都在浪费。而且中间的数据传输链路一旦出问题排查起来非常痛苦。EMR Serverless Spark把这两件事合并到一个平台之后CPU和GPU可以共享同一个弹性资源池。Spark的调度器本身支持GPU资源感知可以在一个作业里同时申请CPU执行器和GPU执行器CPU节点干完特征工程之后数据直接通过内存或分布式缓存交给GPU节点做推理省掉了中间那层网络传输。1.2 异构计算在Serverless环境下的特殊之处在自建Hadoop集群里做异构计算你需要自己装NVIDIA驱动、CUDA toolkit、配置nvidia-docker runtime还要处理Spark和YARN之间的GPU调度逻辑。在EMR Serverless环境里这些底层基础设施都是托管好的你只需要关注两件事一是资源配置写对二是作业代码能识别GPU设备。这里有一个很容易被忽略的点Serverless的弹性伸缩策略和自建集群完全不同。自建集群的GPU节点是常驻的哪怕没任务也在烧钱Serverless是按量计费任务跑完资源立刻释放。这带来的好处是省钱但坏处是如果你在代码里写了“GPU初始化模型”这种耗时操作每次任务启动都要重新加载一次模型这部分时间也在计费。所以在Serverless Spark里做异构计算设计思路要有所调整。我的建议是尽量把模型初始化放在一个长驻的执行器上通过Spark的ReuseExecutionEnvironment机制来复用避免每次作业都重新加载模型。这个细节后面我会在实操环节详细展开。1.3 适用场景与选型判断什么样的作业适合用EMR Serverless Spark做异构计算我总结了三类典型场景第一类是特征工程模型推理的一体化流水线。比如推荐系统里用户特征和物品特征都存在数据湖里每天需要离线更新Embedding向量然后喂给模型做相似度计算。这类任务的特点是数据量大、计算逻辑复杂但模型本身不大百MB级别用单张GPU卡就能跑得动。第二类是批量图片/视频处理。比如内容审核系统需要对海量图片做OCR和分类先用Spark从OSS读取图片元数据再用GPU执行器调用PaddleOCR或者ONNX Runtime做推理。这类任务天然适合Spark的分布式架构因为每张图片的推理是相互独立的。第三类是Spark SQL里嵌入UDF调用GPU算子。如果你有一些计算密集型的自定义函数比如地理空间计算、生物序列比对可以把这部分逻辑放到GPU上执行通过JNI或者JavaCPP调用CUDA库。但不建议用的场景也有如果你的模型特别大比如百亿参数的大模型需要多卡并行或者模型并行那还是老老实实用专门的训练平台吧。Serverless Spark的GPU实例规格有限不适合跑超大模型的训练。还有如果你的ETL逻辑特别简单完全没有AI计算硬塞GPU进去纯属浪费钱。2. 核心配置解析与实操要点2.1 资源规格与CUDA环境适配EMR Serverless Spark的异构配置主要通过Spark配置项来控制核心是这几个参数spark.executor.resource.gpu.amount每个Executor申请多少张GPU卡默认是0需要显式设置为1或更多。spark.task.resource.gpu.amount每个Task占用多少GPU资源通常设置为1.0表示一个Task独占一张卡。spark.executor.resource.gpu.discoveryScriptGPU发现脚本路径。Serverless环境通常会预置但如果你需要自定义可以写一个shell脚本返回GPU设备列表。spark.scheduler.resource.gpu.amount调度队列中GPU资源的配额比例。在实际操作中我建议先把spark.executor.cores和spark.executor.memory的常规资源配好再去加GPU参数。因为如果CPU和内存配置不合理即使GPU配了也会因为执行器启动失败而导致作业报错。注意在EMR Serverless Spark中GPU和CPU并不要求在同一个Executor里。你可以通过配置spark.executor.resource.gpu.amount1让部分Executor只跑GPU任务另一些Executor跑纯CPU任务。这样Spark的调度器会按照资源类型做感知任务会自动路由到有对应资源的节点上。关于CUDA环境坦白说这是Serverless环境省心的地方。阿里云官方镜像里已经预置了CUDA 11.4和cuDNN 8PyTorch和TensorFlow的GPU版本也都装好了。你不需要像自建环境那样纠结驱动版本直接在作业代码里import torch就能用cuda.is_available()检测到设备。2.2 异构计算在Serverless Spark里的三阶段架构从一个可落地的角度来说CPUGPU异构计算在EMR Serverless Spark里可以拆成三个执行阶段这也是我实际项目里打磨出来的标准架构。第一阶段是CPU执行器做数据准备。这段逻辑很简单用Spark SQL或者DataFrame API从数据湖里读数据做过滤、去重、类型转换、特征工程。关键点是最后要把数据转换成模型能直接接受的格式比如固定长度的特征向量、归一化后的数值、Tensor格式的序列。这一步最忌讳的是在CPU端反复序列化和反序列化尽量用mapPartitions在分区内批量处理。第二阶段是GPU执行器做模型推理。具体做法是让GPU执行器去读取一个预设好的模型文件加载到显存里然后接收CPU执行器传来的批量数据。因为Spark的Task是绑定到Executor上的所以只要在配置里指定了GPU资源这个Executor上的Task就能直接访问GPU设备。推理结果可以是一个向量、一个标签或者一组embedding通过collectAsList或者写入分布式缓存返回到CPU端。第三阶段还是CPU执行器做结果汇总。拿到推理结果之后做后处理、阈值判断、写回结果表。理由很简单后处理逻辑往往是IO密集型的比如拼接字符串、关联元数据、写OSS这些交给CPU执行器去并发执行反而更高效GPU的算力应该留给真正的张量运算。2.3 核心参数表下面整理了我验证过的一组基础配置按场景分为全异构CPUGPU都用和纯推理已有数据只需GPU配置项全异构建议值纯推理建议值说明spark.executor.instances42执行器数量按数据量估算spark.executor.cores44每个执行器的CPU核数spark.executor.memory8g8g每个执行器的内存spark.executor.resource.gpu.amount11每个执行器申请1张卡spark.task.resource.gpu.amount1.01.0每个任务占1张卡spark.shuffle.partitions168根据数据量调整不宜过大spark.sql.adaptive.enabledtruetrue开启AQE避免小文件问题spark.sql.adaptive.coalescePartitions.enabledtruetrue动态合并小分区这组配置的原则是CPU执行器就干CPU的活GPU执行器就干GPU的活互不抢资源。如果你在代码里使用了torch.cuda.set_device记得要和Spark分配给这个Task的GPU索引保持一致否则会报“CUDA error: invalid device ordinal”。这个问题我在第四部分会详细说。3. 实操过程与核心环节实现3.1 环境准备与机型选型实操的第一步是先确认EMR Serverless Spark的工作空间支持GPU实例规格。当前主流的GPU实例规格是ecs.gn6i-c4g1.xlarge对应T4卡16G显存和ecs.gn7i-c8g1.2xlarge对应A10卡24G显存。选型逻辑很简单如果你的模型是BERT-base级别约400MBT4足够如果是更大规模的模型或者批量推理并发很高建议直接用A10。在创建项目空间时有几个细节点需要注意存储建议挂载OSS不推荐HDFS。Serverless的Worker是无状态的HDFS只有在EMR集群内才高效OSS配合oss://路径访问性能更稳。队列配置里务必设置资源上限。否则你测试时写错代码导致数据膨胀资源会无限横向扩展账单会很难看。如果同一个项目里要跑CPU和GPU两种作业建议设置两种队列分别绑定不同的资源规格用资源配额做隔离避免CPU作业把GPU队列的资源“借走”。3.2 提交一个GPU推理作业的完整示例我们用一个实际案例来演示从OSS读取100万张图片的元数据用ResNet50模型做批量特征提取结果写回OSS。代码用PySpark编写。from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf from pyspark.sql.types import ArrayType, FloatType import torch import torchvision.models as models from PIL import Image import numpy as np import os # 初始化SparkSession通过conf指定GPU资源 spark SparkSession.builder \ .appName(resnet_feature_extraction) \ .config(spark.executor.resource.gpu.amount, 1) \ .config(spark.task.resource.gpu.amount, 1.0) \ .config(spark.sql.adaptive.enabled, true) \ .getOrCreate() # 在Driver端加载模型用broadcast分发到Executor这里不推荐 # 正确做法是让每个Executor在进程初始化时加载模型 # 这里使用一个全局变量做懒加载 _model None def get_model(): global _model if _model is None: # 根据当前Task分配的GPU索引设置设备 local_rank int(os.environ.get(SPARK_TASK_RESOURCE_GPU_INDEX, 0)) torch.cuda.set_device(local_rank) model models.resnet50(weightsmodels.ResNet50_Weights.IMAGENET1K_V1) model.eval() model.cuda() _model model return _model def extract_feature(image_bytes): import io model get_model() img Image.open(io.BytesIO(image_bytes)).convert(RGB) # 假设已经做了resize和normalize预处理 img_tensor preprocess(img).unsqueeze(0).cuda() with torch.no_grad(): features model(img_tensor).cpu().numpy().flatten() return features.tolist() # 注册UDF并执行推理 extract_feature_udf udf(extract_feature, ArrayType(FloatType())) df spark.read.format(parquet).load(oss://bucket/image_meta/) result_df df.select( col(image_id), extract_feature_udf(col(image_bytes)).alias(features) ) result_df.write.mode(overwrite).parquet(oss://bucket/features/) spark.stop()这里有几个关键细节都是踩过坑之后总结的第一不要在Driver端加载模型后用broadcast发给Executor。模型文件通常几百MBbroadcast传输会占满网络IO而且Executor反序列化模型的时间和直接加载差不多纯属多此一举。第二用全局变量做懒加载是正确的。Executor进程启动后是常驻的_model这个全局变量在进程生命周期里只需要初始化一次。如果每个Task都重新加载模型那作业耗时和费用会成倍增长。第三SPARK_TASK_RESOURCE_GPU_INDEX这个环境变量非常重要。当Executor申请了多张GPU卡时Spark会给每个Task分配一个GPU索引你的代码必须用这个索引来设置CUDA设备不能写死cuda:0。3.3 小批量拼接与显存优化GPU推理最怕的不是算力不够而是显存利用不充分。每张图片单独跑一次模型前向计算显存占用高不说核心利用率可能只有10%。正确的姿势是做batch推理。在上面代码基础上优化思路是用mapPartitions在分区级别做批量处理def batch_extract(iterator): model get_model() images [] ids [] for row in iterator: ids.append(row.image_id) img preprocess(Image.open(io.BytesIO(row.image_bytes))) images.append(img) # 构造一个大batch batch_tensor torch.stack(images).cuda() with torch.no_grad(): features model(batch_tensor).cpu().numpy() for i, feature in enumerate(features): yield (ids[i], feature.tolist())batch size的经验值T4 16G显存跑ResNet50batch size设为32到64比较合适A10 24G显存可以设到64到128。不是越大越好太大了会OOM太小了发挥不了GPU性能。3.4 成本估算与任务监控Serverless的计费逻辑是按量结算CPU和GPU分别计费GPU的单价通常是CPU的5到10倍。这就带来一个很现实的调优目标尽量压缩GPU执行器的空闲时间。怎么判断有没有空闲两个监控指标最关键GPU利用率和显存使用率。EMR Serverless Spark的监控面板里可以按Executor维度查看GPU指标。如果你发现GPU利用率一直在20%以下说明任务瓶颈在CPU端的数据读取或预处理上GPU大部分时间都在等数据。针对这种情况我的优化手段有两种其一调大CPU执行器的并发数让数据准备的速度跟上GPU消费速度其二用spark.sql.adaptive.coalescePartitions.enabled动态合并小分区让每个GPU Task处理更多数据减少Task调度开销。另外一个经验把ETL和推理写在一个作业里但用不同的Stage来区分。通过Spark UI的Stage耗时统计你能很清楚地看到哪些逻辑跑在CPU上、哪些跑在GPU上。如果发现某个Stage耗时异常优先看那个Stage的资源消耗曲线。4. 常见问题与排查技巧实录4.1 典型问题速查表问题现象可能原因排查方法解决方案作业启动失败提示“resource profile not supported”GPU资源规格配置不匹配实例类型不支持GPU检查队列的实例规格更换为GPU规格实例或在队列配置中启用了GPU能力运行时报“CUDA error: invalid device ordinal”代码里写死了cuda:0但Task分配到的GPU索引不是0打印SPARK_TASK_RESOURCE_GPU_INDEX环境变量用环境变量动态设置CUDA设备任务一直Pending排队超时无可用GPU资源或被其他作业占用查看资源队列占用情况调整作业优先级或扩容GPU队列资源上限GPU利用率极低10%数据读取或预处理跟不上GPU在空转查看Stage耗时分布增加CPU执行器数量或优化数据读取格式用parquet替代jsonExecutor进程OOM被Kill显存不够或JVM内存不足查看Executor日志中的OOM信息减小batch size或增大executor.memorySpark执行器启动失败提示“No such file or directory”GPU发现脚本路径不对查看driver日志使用默认脚本路径或修正discoveryScript参数PyTorch无法调用CUDAis_available返回False环境变量CUDA_VISIBLE_DEVICES被清空检查配置在spark.executorEnv里设置CUDA_VISIBLE_DEVICES4.2 资源分配逃逸问题这个坑比较隐蔽我单独拿出来说。在YARN或K8s调度GPU时调度器会把GPU设备“钉死”在容器里。但在Spark的Task调度层面GPU资源是通过环境变量下发给Task的。如果你的Executor申请了多张GPU卡而你的代码没有按照Spark分配的索引去使用GPU就会出现一种情况某个Task的代码即使在执行但它在用另一个Task的设备导致“GPU能跑但性能和显存互相干扰”。解决方法就一句话Task代码里必须读取SPARK_TASK_RESOURCE_GPU_INDEX环境变量用这个值来设置CUDA_VISIBLE_DEVICES。不要自己算设备索引更不要写死。4.3 模型动态加载与版本一致性如果你的模型文件存放在外部存储每次任务启动时从OSS下载要注意版本一致性问题。多个作业并发跑有可能读到模型的中间写入状态导致推理结果异常。建议的解法是模型文件以不可变的方式存放文件路径里带上版本号或commit hash发布新版本就生成新路径不覆盖旧文件。这样即使作业回滚也能准确的找到对应版本的模型。4.4 异构作业级别的容错与重试Spark的Task失败会自动重试但GPU任务重试有个特殊问题如果代码在初始化显存时崩溃显存可能没有完全释放导致同一Executor上重试的Task继续失败。这在低版本CUDA上更容易出现。实操上我建议在代码里加一个保护逻辑当检测到torch.cuda.OutOfMemoryError时先执行torch.cuda.empty_cache()再重新初始化模型。同时把spark.task.maxFailures调小一些比如2不要让同一个Task无限重试避免显存反复泄漏。4.5 跨队列/跨集群依赖的坑另一个容易踩的是数据权限问题。Serverless Spark访问OSS或数据湖使用的是工作空间的执行角色。如果你在同一份流程里CPU作业和GPU作业用的是不同角色GPU执行器读取数据时会遇到权限不足的问题。排查方法是看Executor日志里有没有类似“AccessDenied”报错。如果确认是权限问题最好在项目空间层面统一设置执行角色或者对GPU专用OSS目录单独授权只读权限。个人经验总结我从开始接触EMR Serverless Spark的异构计算到现在大概小半年时间最大的体会是这类Serverless环境最大的优势不在于性能多强而在于把“基础设施的琐碎”全部屏蔽掉了。你不需要关心驱动适配、不需要关心容器运行时、不需要关心GPU调度底层实现专注于业务代码本身就行。但这也带来一个代价传统环境下积累的很多优化经验需要调整。比如在自建GPU服务器上为了省钱你可以手动调整CUDA上下文切换策略但在Serverless环境下你只能在资源配置和代码逻辑上下功夫。最后分享一个实用技巧如果你的模型推理并发量不大但实时性要求高可以把Spark SQL的微批处理改成foreachBatch模式配合流式读入Kafka数据在同一个Spark作业里完成“流式特征计算GPU推理结果写回”这套方案实测能支撑每秒几千次的推理请求而且资源按流量计费闲时几乎不花钱。Serverless Spark的异构能力上限不低值得多花点时间挖掘。
返回列表