You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 22:31:06