PySpark作业出现内存泄漏问题求助
PySpark本地模式作业后内存无法释放的排查与解决
核心问题分析
你遇到的是本地模式(local[*]/local[50])下PySpark作业结束后,JVM内存未被系统正常回收的问题。本地模式中Driver与Executor共享同一个JVM进程,若JVM垃圾回收配置不合理、资源释放不彻底或进程残留,都会导致内存无法释放,必须重启服务器才能恢复。
具体解决方案
1. 统一并优化Spark提交配置
- 消除Master配置冲突:代码中硬编码了
spark.master=local[*],但提交命令又指定--master local[50],两者冲突会导致Spark启动逻辑混乱。删除代码中的.config("spark.master", "local[*]"),统一使用提交命令的参数。 - 调整Driver内存阈值:当前
--driver-memory 200g占用了服务器80%的内存(总251GB),留给系统和JVM GC的空间严重不足。建议下调至180g,给系统预留至少20%的内存用于缓存和GC操作。 - 修正临时目录配置:提交命令中
spark.local.dir=./logs与代码中的/shared/hm31冲突,统一设置为一个有足够空间的目录,且确保该目录不属于tmpfs(避免临时文件占用内存)。
优化后的提交命令示例:
spark-submit --master local[50] \ --driver-class-path ./postgresql-42.5.2.jar \ --jars ./postgresql-42.5.2.jar \ --driver-memory 180g \ --conf "spark.local.dir=/shared/hm31" \ --conf "spark.driver.extraJavaOptions=-XX:+UseG1GC -XX:MaxGCPauseMillis=200" \ calculate_similarities.py
2. 完善代码中的资源释放逻辑
代码中异常分支存在未定义变量的风险(比如result可能还未创建就执行result.unpersist()),导致后续的spark.stop()无法正常执行。同时,unpersist()默认是非阻塞的,需确保数据被彻底释放:
import sys # 封装资源释放逻辑 def clean_resources(spark): if 'result' in locals(): result.unpersist(blocking=True) if 'df' in locals(): df.unpersist(blocking=True) if 'unique_nodes' in locals(): unique_nodes.unpersist(blocking=True) spark.catalog.clearCache() spark.stop() if __name__ == '__main__': spark = SparkSession \ .builder \ .appName("Python Spark SQL basic example") \ .config("spark.driver.extraClassPath", "postgresql-42.5.2.jar") \ .config("spark.executor.extraClassPath","postgresql-42.5.2.jar") \ .config("spark.local.dir", "/shared/hm31") \ .getOrCreate() spark.sparkContext.setLogLevel("WARN") parquet_path = '/shared/hossein_hm31/embeddings_parquets' try: unique_nodes = read_df(spark, jdbc_url, 'hm31.unique_nodes_cert', jdbc_properties) df = spark.read.parquet(parquet_path) unique_nodes.createOrReplaceTempView("unique_nodes") df.createOrReplaceTempView("all_embeddings") sql_query = """ select u.node_id, a.embedding from unique_nodes u inner join all_embeddings a on u.pmid = a.pmid """ result = spark.sql(sql_query) print("num", result.count()) result.repartition(10).write.parquet('/shared/parquets_embeddings/') write_df(result, 'hm31.uncleaned_embeddings_cert', jdbc_properties) clean_resources(spark) sys.exit(0) except Exception as e: print(f'Error: {e}') clean_resources(spark) sys.exit(1)
- 添加
blocking=True确保数据立即从内存中清除 - 通过
locals()检查变量是否存在,避免异常分支报错 - 将资源释放逻辑封装为函数,确保所有分支都能执行到
3. 强制进程退出
在spark.stop()之后,显式触发Python进程退出,避免残留:
sys.exit(0)
若仍有残留,可尝试更强制的退出方式(仅在必要时使用):
import os os._exit(0)
4. 手动清理Spark临时文件
作业结束后,手动清理spark.local.dir下的临时文件:
rm -rf /shared/hm31/*
若该目录是tmpfs挂载的,临时文件会占用内存,必须清理才能释放。
5. 排查残留进程
作业结束后,检查是否有残留的Java/Python进程:
ps aux | grep -E "(java|python)" | grep -v grep
若发现残留进程,直接杀掉:
kill -9 <进程ID>
额外建议
- 避免在本地模式下处理超大规模数据:本地模式不适合处理TB级数据,若数据量持续增长,建议切换到集群模式(Standalone/Yarn/K8s)
- 开启GC日志:在提交命令中添加
-XX:+PrintGCDetails -XX:+PrintGCTimeStamps,分析JVM内存回收情况,定位是否存在真正的内存泄漏
内容的提问来源于stack exchange,提问作者m0ss
相关产品推荐
相关产品推荐

