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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 15:07:20