You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Dataproc集群PySpark作业CPU使用率低及内存错误排查优化

PySpark + SparkNLP 性能问题排查与优化

1. CPU使用率偏低的原因

资源配置不合理

当前Spark配置中spark.executor.cores=1,每个Executor仅使用1核。集群从机为7台c2-standard-8(8核/台),总核数56,但单Executor核数限制导致每个节点要运行8个Executor,调度开销远大于计算开销,无法利用CPU多核优势,整体利用率被拉低。

数据分区失效

加载数据时用created_date_column做repartition键,但每次作业的created_date值完全相同,哈希分区后所有数据集中到少数分区,实际并行度远低于设置的168,大部分Executor处于空闲状态,CPU自然无法跑满。

手动分批的反作用

代码中用limit+subtract实现手动分批,每次循环都会触发全表扫描与shuffle(subtract是宽依赖操作),大量时间消耗在数据准备阶段,真正的NLP计算占比极低,CPU长期处于等待状态。

SparkNLP并行度未充分利用

SentenceDetectorDLModel等组件的默认批次size过小,无法充分压榨CPU算力,导致单任务计算量不足,CPU利用率上不去。

2. 全量解析内存错误的原因

数据倾斜导致单Task内存过载

全量处理时,因created_date值单一,repartition后数据集中在极少数Task中,每个Task要处理远超常规规模的数据量。加上SparkNLP处理时会生成大量Annotator对象(如Document、Sentences),单Task内存占用超过Executor的4G配额,触发OOM。

前置操作的内存消耗

total_rows = remaining_df.select("document_id").count()触发全表扫描,dropDuplicates则需要shuffle所有数据,若分区不合理,Driver或Executor的内存会被大量占用,为后续NLP处理预留的内存不足。

Executor内存配置不足

虽然给每个Executor分配了4G内存,但Spark会预留部分内存用于存储与执行框架,实际可用堆内存有限。SparkNLP的预训练模型本身需要一定内存,大批次数据处理时内存缺口进一步放大,引发bad alloc错误。

3. 性能优化方案

资源配置调整

  • 重新规划Executor资源:将spark.executor.cores设为4,spark.executor.instances设为14(7台从机×2个Executor/台),spark.executor.memory设为12G(每台32G内存扣除系统预留后,分给2个Executor)。此举可减少调度开销,充分利用多核CPU。
  • 调整并行度参数:将spark.default.parallelism和spark.sql.shuffle.partitions设为Executor数量×核数=14×4=56,或按每分区100-200MB的标准调整(1.4GB数据对应10-15个分区即可)。
  • 启用Kryo序列化:在delta_job的Spark配置中添加spark.serializer=org.apache.spark.serializer.KryoSerializer,减少数据序列化开销,降低内存占用。

数据分区优化

  • 更换分区键:摒弃created_date,改用document_id(或其哈希前缀)作为repartition与分区存储的键,保证数据均匀分布到各个分区。
  • 避免过度repartition:若原Delta表已有合理分区,直接加载即可,无需强制repartition到168,根据实际数据量设置分区数。

代码逻辑优化

  • 移除手动分批:Spark原生支持分布式并行处理,无需手动拆分批次。直接全量处理,让Spark自动将任务分配到各个Executor,避免重复扫描与shuffle的开销。
  • 简化去重操作:若document_id为唯一键,将dropDuplicates(["document_id"])改为distinct();若原表已保证数据唯一性,直接去掉该步骤。
  • 优化SparkNLP流水线:
    • 增大模型批次size:给SentenceDetectorDLModel添加setBatchSize(64)(可根据内存情况调整到128),提升模型处理吞吐量。
    • 避免重复fit:流水线中所有组件均为预训练模型或转换器,无需每次调用都fit。修改apply_pipeline方法:
      def apply_pipeline(self, df: DataFrame) -> DataFrame:
          # 仅需fit一次空Schema,避免重复操作
          empty_df = self.spark_session.createDataFrame([], df.schema)
          model = self.pipeline.fit(empty_df)
          return model.transform(df)
      
  • 优化存储逻辑:save_table中无需强制repartition到168,改用document_id作为repartition键,或直接去掉repartition步骤,让Delta Lake自动处理数据分布。

其他建议

  • 关闭动态分配:若集群为专属作业使用,固定Executor数量可避免动态调整带来的开销;若需保留动态分配,设置合理的最小/最大Executor数量。
  • 尽快获取Yarn UI权限:通过UI可查看Executor内存、CPU使用情况,Task运行时长与数据倾斜情况,精准定位剩余性能瓶颈。

内容的提问来源于stack exchange,提问作者Omar Cotugno

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 12:57:03