Dataproc集群PySpark作业CPU使用率低及内存错误排查优化
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)
- 增大模型批次size:给SentenceDetectorDLModel添加
- 优化存储逻辑:
save_table中无需强制repartition到168,改用document_id作为repartition键,或直接去掉repartition步骤,让Delta Lake自动处理数据分布。
其他建议
- 关闭动态分配:若集群为专属作业使用,固定Executor数量可避免动态调整带来的开销;若需保留动态分配,设置合理的最小/最大Executor数量。
- 尽快获取Yarn UI权限:通过UI可查看Executor内存、CPU使用情况,Task运行时长与数据倾斜情况,精准定位剩余性能瓶颈。
内容的提问来源于stack exchange,提问作者Omar Cotugno

