搭配Flask API使用Pyspark时任务完成内存不释放,清缓存无效如何解决?
PySpark + Flask API 服务内存持续上涨问题解决方案
1 易被忽略的内存占用根源
- 仅对最终输出的DataFrame执行
unpersist(),但计算流程中产生的中间RDD/临时DataFrame没有手动释放,很多shuffle、关联操作生成的临时对象如果被隐式缓存,不会随任务结束自动清理 - PySpark的Python进程和JVM进程内存相互独立,你当前的操作仅清理了JVM侧的Spark缓存,Python侧产生的大对象(比如拉取到Driver端的
collect()结果、中间临时Python变量)还残留在Python进程内存中没有释放 - 默认的
df.unpersist()是懒执行操作,不会立即回收内存,需要加参数强制同步清理
2 可直接落地的修复步骤
- 先优化你的缓存清理逻辑,在原有操作基础上补充以下代码,放在每次请求处理完成的位置:
# 强制清理所有未释放的RDD缓存 for (id, rdd) in spark.sparkContext._jsc.getPersistentRDDs().items(): rdd.unpersist(blocking=True) # 显式触发JVM侧垃圾回收 spark.sparkContext._jvm.System.gc()
- 如果你在请求处理中用到了
collect()、toPandas()这类把数据拉取到Driver端的操作,处理完成后要显式删除Python侧的对应大变量,再触发Python GC:
import gc # 替换成你自己的Python侧大对象变量名 del result_pandas_df, temp_data_list gc.collect()
- 启动Spark Session时调整配置,禁用不必要的隐式缓存、优化垃圾回收策略:
spark = SparkSession.builder \ .config("spark.sql.adaptive.enabled", "true") \ .config("spark.sql.adaptive.localShuffleReader.enabled", "false") \ .config("spark.sql.autobroadcastjointhreshold", "-1") # 禁用自动广播,避免广播变量残留 .config("spark.driver.extraJavaOptions", "-XX:+UseG1GC -XX:MaxGCPauseMillis=200") # 给Driver端启用G1垃圾回收器 .getOrCreate()
- 如果你的接口并发量不高,可以放弃单例复用Spark Session的方案,每个请求创建独立的Spark Session,处理完成后直接调用
spark.stop()销毁,虽然有少量Session启动开销,但可以彻底避免跨请求的内存泄漏
3 定位排查方法
- 每次请求前后调用
spark.sparkContext.getExecutorMemoryStatus()、spark.catalog.listTables()查看缓存占用变化,先定位泄漏发生在JVM侧还是Python侧 - 配合
psutil库统计Flask进程的内存占用,结合JVM自带的jstat工具,确认泄漏的具体进程
内容的提问来源于stack exchange,提问作者Jatin Patel - JP
相关产品推荐
相关产品推荐

