Pytest运行PySpark不显示Spark UI且报OOM的解决与调试方案
Spark 3.0 + PySpark pytest 场景OOM与Spark UI不可访问问题解决方案
OOM内存溢出修复
- 优先核对测试环境Spark内存配置:本地pytest运行默认启动的local模式SparkSession,堆内存默认仅分配1G,数据量稍大就会触发OOM。初始化SparkSession时显式指定内存与分区参数,参考配置:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .master("local[4]") \ .appName("pytest_spark_case") \ .config("spark.driver.memory", "4g") \ .config("spark.executor.memory", "4g") \ .config("spark.sql.files.maxPartitionBytes", "134217728") \ # 单读取分区最大128MB,避免单分区过大 .config("spark.sql.adaptive.enabled", "true") \ # 开启Spark 3.0 GA的自适应执行,自动处理分区倾斜、调整shuffle分区大小 .getOrCreate()
注意:local模式下driver与executor共用同一个JVM进程,driver和executor内存参数建议同时配置,避免参数不生效
- 排查文件读取逻辑问题:如果读取的是单个超大文件/未分区的大体积数据集,默认读取后仅生成少量分区,单task承载数据量过大会直接撑爆内存。读取后可显式调用
repartition()调整为合理分区数,保证单分区数据量不超过200MB;读取时添加pathGlobFilter参数过滤目录下无关的大体积文件,避免误读非目标数据。 - 排查pytest资源泄漏与错误写法:不要在每个测试用例中单独初始化SparkSession,使用session级别的pytest fixture全局复用Spark实例,所有用例执行完成后统一调用
stop()回收资源;断言时禁止直接调用toPandas()将全量DataFrame拉取到Python进程,计数类断言直接使用df.count()实现,避免driver端内存被拉取的本地数据占满。fixture参考写法:
import pytest @pytest.fixture(scope="session") def spark(): # 此处补充前述SparkSession初始化逻辑 yield spark spark.stop()
- 补充内存兜底配置:如果堆内存调整后仍有OOM,可开启堆外内存缓解压力,添加配置
.config("spark.memory.offHeap.enabled", "true").config("spark.memory.offHeap.size", "2g")即可。
Spark UI无法访问排查
- 核对Spark UI实际端口:Spark启动时如果4040端口被占用,会自动顺延使用4041、4042等后续端口,不要固定访问4040。启动时查看控制台输出的
Spark UI started at http://<host>:<port>日志,访问日志中标注的实际端口即可。 - 确认UI功能未被手动关闭:如果控制台没有输出UI启动日志,检查Spark配置中是否误设了
spark.ui.enabled=false,显式配置为true即可开启UI。 - 跨网络场景下确认端口连通性:如果在远程服务器/容器内运行pytest,需要确认对应端口没有被防火墙、安全组拦截,通过本地端口转发映射后再访问;如果是CI流水线环境无直接访问条件,直接使用下述无UI调试方案即可,无需强行访问实时UI。
无Spark UI场景下的PySpark作业调试方法
- 开启Spark事件日志留存全量诊断数据:初始化SparkSession时添加配置,将所有task执行指标、shuffle数据量、异常栈信息持久化到本地磁盘,作业运行结束后可随时通过Spark History Server加载查看,不依赖运行时的实时UI。配置参考:
.config("spark.eventLog.enabled", "true") .config("spark.eventLog.dir", "file:///tmp/spark-events") # 需提前手动创建该目录
- 打印执行计划提前排查逻辑问题:在触发action操作(比如
count())之前,调用df.explain(mode="formatted")打印结构化执行计划,排查是否存在意外的笛卡尔积、分区裁剪失效、隐式类型转换导致过滤不生效等问题,这类逻辑问题通常会导致全量数据被拉取到单点触发OOM。 - 调整日志级别保留关键运行信息:不要为了减少控制台输出把Spark日志级别设为ERROR,调整为WARN级别即可在控制台看到内存不足预警、数据倾斜提示、task失败原因等关键信息,大部分问题直接看控制台日志就能定位,不需要进UI看指标。
- 小样本递进式验证逻辑:跑全量数据前先取小样本验证流程,比如先执行
df.limit(1000).count()确认逻辑可跑通,再逐步放大数据量,快速定位是数据量超过内存阈值导致的OOM,还是逻辑本身存在问题。 - 内置指标快速定位倾斜:如果怀疑是数据倾斜导致的OOM,可以在count之前遍历各分区的记录数,查看是否存在单分区数据量远高于均值的情况,示例代码:
from pyspark.sql.functions import spark_partition_id partition_counts = df.groupBy(spark_partition_id()).count().collect() print(partition_counts)
内容的提问来源于stack exchange,提问作者Xi12
相关产品推荐
相关产品推荐

