PySpark的Py4J错误:为何一段脚本运行正常另一段报错?
PySpark运行差异与配置疑问
问题场景
- 以下代码可正常运行:
spark = SparkSession.builder.appName('pyspark').getOrCreate() book_local = spark.read.text("data.txt") book_local.show()
- 运行以下代码时抛出错误:
spark = SparkSession.builder.appName('pyspark').getOrCreate() my_grocery_list = [ ["Banana", 2, 1.74], ["Apple", 4, 2.04], ["Carrot", 1, 1.09], ["Cake", 1, 10.99], ] df_grocery_list = spark.createDataFrame(my_grocery_list) df_grocery_list.show() # 此处抛出错误
- 错误信息:
Py4JJavaError: java.io.IOException: Cannot run program "python3"
- 设置环境变量后恢复正常:
import os import sys from pyspark.sql import SparkSession os.environ['PYSPARK_PYTHON'] = sys.executable os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable
用户疑问
- 为何第一段代码运行无问题,第二段却抛出Py4J错误?第一段代码完全不使用Py4J包吗?
- 尝试用以下代码替换环境变量配置时仍报错,原因是什么?
spark = SparkSession.builder.appName('pyspark').config("spark.pyspark.python", sys.executable).getOrCreate()
解答
两段代码的运行差异原因
第一段代码不是完全不用Py4J,只是执行逻辑没触发Spark启动Python工作节点:
spark.read.text("data.txt")读取文本文件时,Spark底层用Java/Scala原生组件处理文件读取,show()展示结果只是通过Py4J把JVM端的结果传回Python端,全程不需要启动独立的Python工作进程,也就不会触发python3命令的调用。- 第二段代码里,
spark.createDataFrame(my_grocery_list)需要把本地Python集合的数据序列化传给JVM,后续show()触发执行时,Spark要启动Python工作节点来完成数据转换和计算,这时候会尝试调用python3命令启动进程,如果系统环境没配置好python3的正确路径,就会抛出无法运行的错误。
config("spark.pyspark.python")配置无效的原因
用builder.config()设置参数无效,核心是SparkSession的配置必须在SparkContext初始化前设置:
- 如果之前运行过代码,
getOrCreate()会复用已经存在的SparkSession,这时候后续添加的config参数不会生效。 - 解决方法是先停止已有的会话,再重新构建:
import sys from pyspark.sql import SparkSession # 先停止现有会话 try: spark.stop() except: pass # 重新构建并配置参数 spark = SparkSession.builder.appName('pyspark')\ .config("spark.pyspark.python", sys.executable)\ .getOrCreate()
内容的提问来源于stack exchange,提问作者telecomshy
相关产品推荐
相关产品推荐

