PySpark 3.5.0调用show()时Python Worker意外退出报错求助
PySpark 3.5.0 运行
createDataFrame时出现Python Worker崩溃错误 我刚安装了PySpark 3.5.0,运行以下代码:
data = [("Java", "20000"), ("Python", "100000"), ("Scala", "3000")] df = spark.createDataFrame(data) df.show()
触发如下错误:
24/02/19 11:41:39 ERROR Executor: Exception in task 0.0 in stage 0.0 (TID 0)/ 1] org.apache.spark.SparkException: Python worker exited unexpectedly (crashed) at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:612) at org.apache.spark.api.python.BasePythonRunner$ReaderIterator$$anonfun$1.applyOrElse(PythonRunner.scala:594) at scala.runtime.AbstractPartialFunction.apply(AbstractPartialFunction.scala:38) at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:789) at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:766) 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 scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460) at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source) at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43) at org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43) at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:890) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:890) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) 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.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.io.EOFException at java.base/java.io.DataInputStream.readInt(DataInputStream.java:397) at org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:774) ... 26 more 24/02/19 11:41:39 WARN TaskSetManager: Lost task 0.0 in stage 0.0 (TID 0)
试过Anaconda和pip两种方式安装PySpark,错误完全一致,求可行的解决方法。
常见解决方向
- 验证Python版本兼容性:PySpark 3.5.0仅支持Python 3.8~3.11版本,用
python --version或python3 --version确认当前版本,不符合则切换到兼容版本。 - 指定PySpark使用的Python解释器:确保Spark调用的Python和安装PySpark的环境一致,可通过设置环境变量解决:
终端执行代码前先运行:
或者在代码开头添加:export PYSPARK_PYTHON=/your/python/path export PYSPARK_DRIVER_PYTHON=/your/python/pathimport os os.environ['PYSPARK_PYTHON'] = '/your/python/path' os.environ['PYSPARK_DRIVER_PYTHON'] = '/your/python/path' - 检查Java版本:PySpark 3.5.0要求Java 8或Java 11,Java 17存在兼容性问题,用
java -version确认版本,不符合则切换。 - 清理重装PySpark:
- 卸载现有PySpark:
pip uninstall pyspark -y(conda环境用conda remove pyspark -y) - 清理缓存:
pip cache purge(conda环境用conda clean --all) - 重新安装指定版本:
pip install pyspark==3.5.0
- 卸载现有PySpark:
- 禁用动态资源分配:初始化SparkSession时添加配置:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("test") \ .config("spark.dynamicAllocation.enabled", "false") \ .getOrCreate() - 查看Python Worker崩溃日志:Spark会在临时目录生成Worker日志,通过
spark.sparkContext.getConf().get("spark.local.dir")获取路径,查看日志定位具体错误(如依赖缺失、权限问题等)。
内容的提问来源于stack exchange,提问作者Mihai Tache
相关产品推荐
相关产品推荐

