pyspark3.1异常: Python worker failed to connect back

pyspark环境配置报错解决

异常描述

环境:win10, spark3.1.2版本,hadoop3.3.1,java1.8 在pycharm或直接在pyspark shell环境中执行如下测试代码报错: pyspark3.1: Python worker failed to connect back

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, LongType, StringType, IntegerType

if __name__ == "__main__":
   spark = SparkSession.builder.master(local[1]).getOrCreate()

   spark_rdd = spark.sparkContext.parallelize([
    (123, "Katie", 19, "brown"),
    (456, "Michael", 22, "green"),
    (789, "Simone", 23, "blue")])

   # 设置dataFrame将要使用的数据模型,定义列名,类型和是否为能为空
   schema = StructType([StructField("id", IntegerType(), True),
                        StructField("name", StringType(), True),
                        StructField("age", IntegerType(), True),
                        StructField("eyeColor", StringType(), True)])
   # 创建DataFrame
   spark_df_from_rdd = spark.createDataFrame(spark_rdd, schema)
   spark_df_from_rdd.show()

异常信息:

23:16:58 ERROR Executor: Exception in task 0.0 in stage 0.0 (TID 0)
org.apache.spark.SparkException: Python worker failed to connect back.
        at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:170)
        at org.apache.spark.api.python.PythonWorkerFactory.create(PythonWorkerFactory.scala:97)
        at org.apache.spark.SparkEnv.createPythonWorker(SparkEnv.scala:117)
        at org.apache.spark.api.python.BasePythonRunner.compute(PythonRunner.scala:108)
        at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:324)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:288)
        at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
        at org.apache.spark.scheduler.Task.run(Task.scala:121)
        at org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:402)
        at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:408)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
        at java.lang.Thread.run(Unknown Source)
Caused by: java.net.SocketTimeoutException: Accept timed out

解决方法

win10增加系统环境变量:

key: PYSPARK_PYTHON
 value: python

问题成功解决!!!

经验分享 程序员 微信小程序 职场和发展