设置spark.executor.pyspark.memory后PySpark集群遇导入错误求助
我在AWS SageMaker Processing环境(基于Docker构建,包含所有代码与pip依赖,采用全局pip安装Python)中运行的PySpark集群遭遇OOM错误。为避免Python进程占用JVM内存,尝试配置spark.executor.pyspark.memory参数,但设置后出现Pandas相关导入错误,移除该配置则错误消失。请问是否需要调整其他Spark配置、Python/pip设置,或处理spark --py-files/源库复制?
错误信息
23:02:10.654 [task-result-getter-0] WARN org.apache.spark.scheduler.TaskSetManager - Lost task 3.0 in stage 17.0 (TID 237) (algo-19 executor 15): org.apache.spark.api.python.PythonException: Traceback (most recent call last): File "/usr/local/lib/python3.9/dist-packages/pandas/__init__.py", line 22, in <module> from pandas.compat import is_numpy_dev as _is_numpy_dev File "/usr/local/lib/python3.9/dist-packages/pandas/compat/__init__.py", line 15, in <module> from pandas.compat.numpy import ( File "/usr/local/lib/python3.9/dist-packages/pandas/compat/numpy/__init__.py", line 4, in <module> from pandas.util.version import Version File "/usr/local/lib/python3.9/dist-packages/pandas/util/__init__.py", line 1, in <module> from pandas.util._decorators import ( # noqa:F401 File "/usr/local/lib/python3.9/dist-packages/pandas/util/_decorators.py", line 14, in <module> from pandas._libs.properties import cache_readonly # noqa:F401 File "/usr/local/lib/python3.9/dist-packages/pandas/_libs/__init__.py", line 13, in <module> from pandas._libs.interval import Interval ImportError: /usr/local/lib/python3.9/dist-packages/pandas/_libs/interval.cpython-39-x86_64-linux-gnu.so: failed to map segment from shared object at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.handlePythonException(PythonRunner.scala:572) at org.apache.spark.sql.execution.python.PythonArrowOutput$$anon$1.read(PythonArrowOutput.scala:118) at org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:525) at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:491) at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460) at org.apache.spark.sql.execution.columnar.DefaultCachedBatchSerializer$$anon$1.hasNext(InMemoryRelation.scala:119) at org.apache.spark.sql.execution.columnar.CachedRDDBuilder$$anon$2.hasNext(InMemoryRelation.scala:286) at org.apache.spark.storage.memory.MemoryStore.putIterator(MemoryStore.scala:223) at org.apache.spark.storage.memory.MemoryStore.putIteratorAsValues(MemoryStore.scala:302) at org.apache.spark.storage.BlockManager.$anonfun$doPutIterator$1(BlockManager.scala:1601) at org.apache.spark.storage.BlockManager.org$apache$spark$storage$BlockManager$$doPut(BlockManager.scala:1528) at org.apache.spark.storage.BlockManager.doPutIterator(BlockManager.scala:1592) at org.apache.spark.storage.BlockManager.getOrElseUpdate(BlockManager.scala:1389) at org.apache.spark.storage.BlockManager.getOrElseUpdateRDDBlock(BlockManager.scala:1343) at org.apache.spark.rdd.RDD.getOrCompute(RDD.scala:376) at org.apache.spark.rdd.RDD.iterator(RDD.scala:326) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750)
Spark配置(spark-defaults.conf)
spark.driver.memory=191284m spark.yarn.am.cores=48 spark.executor.memory=103293m spark.executor.cores=47 spark.executor.instances=19 spark.driver.maxResultSize=19128m spark.executor.memoryOverhead=34431m spark.executor.pyspark.memory=34431m spark.sql.adaptive.enabled=true spark.sql.adaptive.skewJoin.enabled=true spark.default.parallelism=38 spark.task.cpus=23 spark.sql.execution.arrow.pyspark.enabled=true spark.sql.execution.arrow.pyspark.fallback.enabled=true spark.network.timeout=600s spark.serializer=org.apache.spark.serializer.KryoSerializer spark.kryoserializer.buffer.max=1g spark.driver.extraJavaOptions=-XX:+PrintGCDetails -XX:+PrintGCTimeStamps spark.executor.extraJavaOptions=-XX:+PrintGCDetails -XX:+PrintGCTimeStamps spark.eventLog.enabled=true spark.eventLog.dir=/opt/ml/processing/output Executor memory: 172155
问题分析与解决方案
核心原因
设置spark.executor.pyspark.memory后,Python进程的内存被严格限制,导致Pandas加载底层C扩展库(.so文件)时内存不足,无法完成内存映射。这个错误本质是Python进程的可用内存不足以支撑Pandas库的初始化。
调整建议
合理分配Python内存配额
当前spark.executor.pyspark.memory设置为34431m,占spark.executor.memory的1/3左右,但可能挤占了系统或JVM的可用内存。建议降低该值至spark.executor.memory的1/4(约25000m),同时确保JVM有足够内存运行任务。注意该参数是从spark.executor.memory中划分的,并非额外内存。优化Overhead与系统内存
检查spark.executor.memoryOverhead配置,当前34431m可适当增加(比如到40000m),给系统预留更多内存用于加载共享库。同时确认容器/实例的物理内存能支撑executor memory + memoryOverhead的总消耗。精简Python环境
清理Docker镜像中冗余的pip依赖,只保留任务必需的库,减少内存占用。重新构建镜像时指定稳定版本的Pandas(如pip install pandas==1.5.3),避免最新版本的兼容性问题。调整任务并行度
当前spark.task.cpus=23,每个executor仅能运行2个任务,可降低该值(如设为10),增加单executor的任务数,分散内存压力。同时检查数据分区,确保单个任务处理的数据量不会过大。优化共享库加载
在Docker镜像中添加环境变量MALLOC_ARENA_MAX=4,限制内存分配器的arena数量,减少内存碎片,帮助Pandas顺利加载底层扩展库。
内容的提问来源于stack exchange,提问作者Nick Anderson

