PySpark创建DataSet失败,加载DataFrame正常,疑与JAVA_HOME有关
Spark创建DataSet崩溃问题排查
问题描述
刚开始学习Scala,遇到疑似与JAVA_HOME环境变量相关的问题:加载CSV生成DataFrame的代码可正常运行,但创建新DataSet时程序崩溃。
正常运行代码
# Aqui hacemos lo mismo añadiendo la aplicación spark2 = (SparkSession.builder .master("local") .config('spark.ui.port','4050') .getOrCreate()) # Reading csv into a variable to store it datos = spark2.read.csv('housing.csv',header = True).show()
崩溃代码
df = spark.createDataFrame([ Row(a=1, b=2., c='string1', d=date(2000, 1, 1), e=datetime(2000, 1, 1, 12, 0)), Row(a=2, b=3., c='string2', d=date(2000, 2, 1), e=datetime(2000, 1, 2, 12, 0)), Row(a=4, b=5., c='string3', d=date(2000, 3, 1), e=datetime(2000, 1, 3, 12, 0)) ]) df.show()
错误信息
Py4JJavaError: An error occurred while calling o76.showString. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1.0 failed 1 times, most recent failure: Lost task 0.0 in stage 1.0 (TID 1) (192.168.1.13 executor driver): org.apache.spark.SparkException: Python worker failed to connect back.
解决方案
针对Python worker failed to connect back错误,可从以下几个方向排查:
- 检查JAVA_HOME配置:确保系统环境变量中JAVA_HOME指向正确的JDK路径(Spark推荐使用JDK 8或11),且Spark的
spark-env.sh(Linux)或spark-env.cmd(Windows)中也配置了一致的JAVA_HOME,避免Spark启动时找不到正确的Java环境。 - 验证Python与Spark版本兼容性:确认使用的Python版本与Spark版本匹配(例如Spark 3.0+支持Python 3.6及以上),版本不兼容会导致Python worker初始化失败。
- 补全代码依赖导入:崩溃代码中使用了
Row、date、datetime,需在代码开头添加导入语句:
缺少导入会导致worker执行时出现找不到类的错误,进而引发连接失败。from pyspark.sql import Row from datetime import date, datetime - 指定Spark使用的Python解释器:在创建SparkSession时显式指定Python路径,避免环境不一致:
spark2 = (SparkSession.builder .master("local") .config('spark.ui.port','4050') .config("spark.pyspark.python", "/usr/bin/python3") # 替换为你的Python路径 .getOrCreate()) - 检查端口与防火墙:Spark driver与worker通过本地端口通信,若端口被占用或防火墙拦截,会导致连接失败。可尝试关闭本地防火墙,或通过
spark.driver.port配置指定未被占用的端口。
内容的提问来源于stack exchange,提问作者Luis Gómez-Jordana Martín
相关产品推荐
相关产品推荐

