PySpark首次运行df.show()正常,二次运行失败求助
Spark二次运行DataFrame代码后df.show()报错:Python worker无法回连、Socket超时
问题现象
首次运行DataFrame创建代码时,df.show()可正常输出数据;未修改代码二次运行时,df.show()执行失败,抛出Py4JJavaError,核心错误为Python worker无法回连及Socket连接超时。已尝试清除内核缓存、重新上传文件,问题未解决。
相关代码
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate()
from datetime import datetime, date import pandas as pd from pyspark.sql import Row 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 df = spark.createDataFrame([ (1, 2., 'string1', date(2000, 1, 1), datetime(2000, 1, 1, 12, 0)), (2, 3., 'string2', date(2000, 2, 1), datetime(2000, 1, 2, 12, 0)), (3, 4., 'string3', date(2000, 3, 1), datetime(2000, 1, 3, 12, 0)) ], schema='a long, b double, c string, d date, e timestamp') df pandas_df = pd.DataFrame({ 'a': [1, 2, 3], 'b': [2., 3., 4.], 'c': ['string1', 'string2', 'string3'], 'd': [date(2000, 1, 1), date(2000, 2, 1), date(2000, 3, 1)], 'e': [datetime(2000, 1, 1, 12, 0), datetime(2000, 1, 2, 12, 0), datetime(2000, 1, 3, 12, 0)] }) df = spark.createDataFrame(pandas_df) df # All DataFrames above result same. df.show() df.printSchema()
报错信息
Py4JJavaError: An error occurred while calling o421.showString. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 0.0 failed 1 times, most recent failure: Lost task 0.0 in stage 0.0 (TID 0) (PainPeko executor driver): org.apache.spark.SparkException: Python worker failed to connect back. ...(完整报错堆栈信息)
解决方案
1. 显式停止并重建SparkSession
每次运行代码前,先停止已存在的SparkSession,避免残留的Python worker进程干扰:
from pyspark.sql import SparkSession # 停止已有Session(忽略不存在的情况) try: spark.stop() except: pass spark = SparkSession.builder.getOrCreate()
2. 调整Spark Python Worker配置
在初始化SparkSession时添加参数,禁用worker复用并增大超时时间:
spark = SparkSession.builder \ .config("spark.python.worker.reuse", "false") \ .config("spark.network.timeout", "300s") \ .config("spark.executor.heartbeatInterval", "60s") \ .getOrCreate()
spark.python.worker.reuse: 禁用worker进程复用,每次任务启动新进程,避免残留进程导致的连接问题spark.network.timeout: 延长网络超时时间,防止Socket连接超时spark.executor.heartbeatInterval: 调整executor心跳间隔,降低心跳超时概率
3. 确保Python环境一致性
保证Driver和Executor使用的Python版本完全一致,可通过配置指定Python路径:
spark = SparkSession.builder \ .config("spark.pyspark.python", "/usr/bin/python3") \ .config("spark.pyspark.driver.python", "/usr/bin/python3") \ .getOrCreate()
替换路径为你环境中实际的Python可执行文件路径。
4. 手动清理残留进程
在Linux/macOS环境下,清理残留的Spark Python worker进程:
# 查找并杀死pyspark相关进程 ps aux | grep pyspark | grep -v grep | awk '{print $2}' | xargs kill -9
Windows环境可通过任务管理器找到并结束相关Python进程。
内容的提问来源于stack exchange,提问作者Window Man
相关产品推荐
相关产品推荐

