PySpark多转换步骤内存溢出与重复计算问题咨询
核心结论
你修改为复用同一变量的方案不能解决上述两个问题,Spark DataFrame为不可变结构,无论赋值给新变量还是覆盖原有变量,都会生成新的血缘链路,不会改变实际执行逻辑。以下是针对两个问题的具体解决方案:
1. 重复计算问题解决
Spark所有转换算子为懒执行模式,每次触发show/write/count等action操作时,都会从头执行完整的血缘链路,你的耗时叠加本质是每次都重复执行了相似度计算UDF。
解决方法为对UDF计算完成后的中间结果做持久化:
from pyspark import StorageLevel # UDF计算完成后持久化,优先存内存,内存不够放磁盘 df_new_2 = df_spk.select("Orig_doc_id", func_udf(col("text")).alias('Output')) df_new_2.persist(StorageLevel.MEMORY_AND_DISK) # 第一次触发action后,后续操作直接读取缓存,无需重跑UDF df_new_2.show() df_new_3 = df_new_2.select("Orig_doc_id","Output.*") df_new_3.show() # 不会再重复计算UDF,耗时大幅降低
2. YARN内存超限问题解决
报错原因是PySpark Python UDF的内存占用属于堆外内存,默认分配的spark.yarn.executor.memoryOverhead配额不足:
- 首先提交作业时调大堆外内存配额,比如增加配置:
--conf spark.yarn.executor.memoryOverhead=8192,给executor分配8G堆外内存 - 优化UDF内存占用:你用于相似度匹配的语料库如果是加载在
func_dict中,需要将语料库封装为广播变量分发到各节点,避免每个task都复制一份语料库副本,大幅降低内存消耗 - 优化Join逻辑:如果
df_spk_train是小表,使用广播Join避免shuffle:
from pyspark.sql.functions import broadcast df_new_5 = df_new_4.join(broadcast(df_spk_train), df_new_4.max_sim_origin == df_spk_train.text,"inner")
另外你修改后的复用变量代码存在语法错误,filter步骤引用了已经不存在的df_new_3变量,无需修改原有多变量赋值逻辑,变量名仅为DataFrame的引用标识,不会带来额外性能损耗。
内容的提问来源于stack exchange,提问作者Aparajith Chandran
相关产品推荐
相关产品推荐

