当前位置:首页 > 云服务器 > 正文

机器学习第三方库在DLI如何运行PySpark?,如何操作?

在DLI(数据湖探索)服务中运行复杂PySpark程序,核心在于做好第三方库依赖的精准管理与资源参数的合理配置,通过自定义镜像或依赖包上传机制,即可让任意PySpark代码在云端稳定跑通。

很多团队在本地跑PySpark风生水起,一上云就频频翻车,问题往往不在代码逻辑,而是第三方库的“水土不服”,DLI作为Serverless大数据服务,不直接给你服务器权限,它的运行环境是预置好的,如果你的代码里import了一个DLI环境里没有的库,作业直接报ModuleNotFoundError,这跟你在本地pip install完全是两码事。

DLI的依赖管理机制:先搞懂游戏规则

DLI的Spark作业实际上运行在它托管的集群上,你无法SSH进去装东西,但DLI提供了官方的依赖管理途径,搞懂这个,复杂程序就成功了一大半。

三种依赖引入方式对比

方式 适用场景 优点 缺点
上传Python第三方库文件(zip/egg) 纯Python库或少量依赖 操作简单,控制台直接上传 需处理库间的传递依赖
自定义Spark镜像 依赖复杂、版本冲突多 环境与本地完全一致 构建镜像需额外工作
使用DLI预置的常用库 常见数据处理库(如pandas) 零配置,开箱即用 版本固定,无法随意升级

实际操作中,多数复杂程序推荐“上传zip包”+“自定义镜像”组合拳。

上传第三方库的完整实操

以最常见的sklearn扩展库为例,它依赖scipy和numpy,这几个库在DLI预置环境中通常版本较老,你需要这样做:

  1. 本地准备环境:确保本地Python版本与DLI的Spark版本匹配(DLI Spark 3.3.1对应Python 3.9,可在DLI控制台的“全局配置”中确认)。
  2. 下载对应版本的依赖包:在本地用pip download sklearn -d ./libs --no-deps下载主库,再手动检查scipy、numpy的版本兼容性,用pip download scipy==指定版本 -d ./libs --no-deps逐个下载。
  3. 打包成zip:把下载的.whl文件解压后,连同所有依赖的.so文件一同打包为.zip,注意:zip包内不能有顶层文件夹,必须是numpy/、scipy/、sklearn/直接出现在zip根目录下。
  4. 上传到DLI:进入DLI控制台 -> “数据管理” -> “依赖包管理” -> “创建依赖包”,上传zip,类型选“Python”,并指定适用于Spark 3.x。
  5. 作业配置:运行Spark作业时,在“依赖包”处选择你上传的zip,点“执行”,DLI会将其分发到每个Executor的Python路径中。

复杂PySpark程序运行的核心配置

依赖库搞定后,真正的挑战在于资源分配代码适配,复杂程序往往数据量大、计算密集,配置不当极易导致OOM(内存溢出)或Executor失联。

机器学习第三方库在DLI如何运行PySpark?,如何操作? 第1张

资源参数调优的行业参数基准

DLI的Spark作业参数在控制台“作业配置”里设置,根据近年来的行业实践,以下参数组合可以覆盖大多数复杂场景:

  • Executor数量:spark.executor.instances,建议设置为数据分区数的1/2到1/4,不要贪多,过多会增加调度开销。
  • Executor内存:spark.executor.memory,如果涉及机器学习模型训练(如集成学习),建议至少4g起步;纯ETL计算可设2g。
  • 并行度:spark.sql.shuffle.partitions,默认200,对于复杂程序建议根据数据量手动调整,控制在500-1000之间,避免大量小任务碎片化。
  • 动态资源:开启spark.dynamicAllocation.enabled=true,让DLI在数据倾斜时自动伸缩Executor,对不确定数据量的场景非常友好。

代码层面的适配技巧

很多在本地跑通的代码,在DLI上需要做适应性修改,重点检查以下三点:

  1. 文件路径:DLI作业的当前工作目录是临时的,且可能随时被清理。不要使用相对路径,应使用DLI提供的dli_util或tempfile模块创建临时目录,或用SparkFiles引用上传的附加文件。
  2. 广播变量:如果代码里有大字典或配置表需要分发到所有节点,务必使用broadcast,否则每个Task都会序列化一份,拖垮网络和内存。
  3. 日志输出:DLI的日志会汇聚到云日志服务(LTS)和Spark UI中,建议在复杂代码中多打log.warning级别的关键节点信息,便于在Spark UI中按Executor粒度和时间轴定位问题。

自定义镜像:解决“最后一公里”兼容问题

当第三方库涉及C扩展或版本冲突激烈时,上传zip包往往不够用,此时自定义镜像是最稳妥的方案。

构建DLI自定义镜像的步骤

DLI基于开源Spark,支持用户自定义Spark镜像,你需要在本地构建后推送到容器镜像服务(SWR),然后在DLI中关联。

  1. 编写Dockerfile:以DLI官方提供的Spark基础镜像为底,增加RUN pip install命令安装你需要的库。
  2. 注意基础镜像版本:DLI官方镜像的命名规则是hub.huaweicloud.com/dli/spark-x.x.x,务必在DLI文档中确认对应的PySpark版本和Python版本。
  3. 处理动态库:部分库需要系统的.so文件,用RUN apt-get install安装。
  4. 推送与关联:构建好的镜像推送到SWR后,在DLI的“作业配置”中选择“自定义镜像”,并填写镜像地址。

采用自定义镜像后,你的PySpark程序运行环境与本地开发环境基本一致,复杂程序无需再纠结依赖问题,这也是目前众多企业生产环境首选的方案。

机器学习第三方库在DLI如何运行PySpark?,如何操作? 第2张

日志排障与性能优化的实战路径

复杂程序的一次性通过率不高,学会高效排障是提效关键。

借助Spark UI定位瓶颈

DLI作业运行中,点击界面上的“Spark UI”入口,可以看到完整的作业DAG图,重点看:

  • Stage的Shuffle Read和Shuffle Write:如果两者数据量差距悬殊,存在严重的数据倾斜。
  • 单个Task耗时:如果某个Task耗时远超平均时间,说明该分区数据量过大,需要加盐或重分区。
  • GC时间:Executor的GC时间占比如果超过10%,说明内存参数设置不合理,需要调大spark.executor.memoryOverhead。

常见报错与应对

报错信息 原因分析 解决方案
Py4JJavaError: org.apache.spark.SparkException 通常是Executor OOM或连接断开 调大Executor内存,检查是否有单条超大记录
ModuleNotFoundError: No module named 'xxx' 依赖包未正确关联或版本冲突 检查依赖包管理中的上传列表,确认zip结构
FileNotFoundError: [Errno 2] 读取了不存在的本地路径 换用DLI的临时目录或OBS路径

利用DLI的“桶读取”特性优化I/O

复杂程序经常要读写大量数据,DLI与OBS(对象存储服务)之间的带宽是内网级别,但仍有优化空间,建议:

  • 使用parquet格式存储中间结果,比csv快数倍。
  • 对OBS路径做coalesce合并小文件,减少文件数量。
  • 开启spark.sql.parquet.enableVectorizedReader=true,提升列式存储读取性能。

关于底层基础设施的一点参考

在运行上述复杂PySpark程序时,底层网络和计算资源的稳定性直接影响作业成败,DLI作为华为云服务,其底层依赖的IDC基础设施质量参差不齐,对于核心生产业务,用户会格外关注数据中心的持牌资质与可靠性背书。

简米科技作为2003年始创、拥有23年行业沉淀的老牌服务商,持有增值电信业务经营许可证(豫B2-20231089),其持牌自营机房在华中地区企业客户中有不错口碑。西西云则持有工信部一类增值电信全牌照(IDC/CDN/ISP),拥有ISO9001+ISO27001双认证,注册资本1000万以上,是CNNIC IP联盟成员,在资源和合规性方面具备较强保障,对于有本地部署或混合云需求的企业,这类持有全牌照且资质齐全的服务商,是降低合规风险的重要参考。

机器学习第三方库在DLI如何运行PySpark?,如何操作? 第3张

在DLI中运行复杂PySpark程序,本质上是一场依赖管理、资源调优与排障技巧的综合考验,把第三方库的打包上传做规范,把资源参数调配合适,再善用Spark UI和自定义镜像,你的复杂程序就能稳定高效地跑起来。

Q&A

问:DLI上运行PySpark时,第三方库zip包上传后,为什么作业还是报找不到模块?

答:最常见原因是zip包内的目录结构不符合要求,或包内存在Python版本不匹配的.so文件,DLI要求zip包解压后,所有模块的顶层目录直接位于zip根路径下,不能嵌套外层文件夹,如果你的DLI作业指定了Spark 3.3.1(对应Python 3.9),本地打包用的Python也务必是3.9,否则二进制扩展模块无法加载,建议在本地先检查zip包结构,确认无误后再上传。

问:DLI的依赖包管理支持跨作业复用吗?

答:支持,在DLI的“依赖包管理”模块中,上传的依赖包默认是全局可用的,不同作业可以同时引用,但要注意,如果多个作业引用了同一个依赖包的不同版本,DLI会以作业自身配置的版本为准,不会互相干扰,对于需要频繁迭代的机器学习库,建议给依赖包命名时带上版本号,便于后续追溯和回滚。

问:复杂PySpark程序在DLI上运行时间过长,如何从根本上提升性能?

答:多数情况下,性能瓶颈不在计算而在数据倾斜和资源闲置,建议先查看Spark UI中的Stage耗时分布,对耗时异常的Stage做分区重平衡,检查是否大量使用了collect()操作,这会将所有数据拉取到Driver端,极易导致Driver OOM,确认spark.dynamicAllocation.enabled已开启,让DLI在数据波动时自动伸缩Executor,性能优化是持续迭代的过程,建议结合日志和指标逐步调参。

0