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

PySpark DataFrame.show()报错排查及JSON转Spark DataFrame高效方法咨询

Spark DataFrame.show()报错解决及高效JSON转Spark DataFrame方案

问题背景

我是Python新手,已成功创建Spark DataFrame sparkDF,但调用.show()方法时触发报错。当前采用先将JSON转为Pandas DataFrame再转为Spark DataFrame的流程,现寻求更高效的JSON转Spark DataFrame方法。

原代码

!pip install pyspark
import pyspark
from pyspark.sql import SparkSession
import pandas as pd
import json

spark = SparkSession.builder.appName("SparkTrial").config("spark.some.config.option", "some-value").getOrCreate()

f = open('data.json')
data = json.load(f)

pddf=pd.json_normalize(data, "results")

sparkDF = spark.createDataFrame(pddf)

print(sparkDF) # 输出: DataFrame[v: double, vw: double, o: double, c: double, h: double, l: double, t: bigint, n: bigint]

sparkDF.show() # 报错发生在这一行

报错信息(中文翻译)

---------------------------------------------------------------------------
Py4JJavaError                             Traceback (most recent call last)
输入 In [7], in <cell line: 1>()
----> 1 sparkDF.show()

文件 ~\anaconda3\lib\site-packages\pyspark\sql\dataframe.py:606, in DataFrame.show(self, n, truncate, vertical)
   603     raise TypeError("Parameter 'vertical' must be a bool")
   605 if isinstance(truncate, bool) and truncate:
---> 606     print(self._jdf.showString(n, 20, vertical))
   607 else:
   608     try:

文件 ~\anaconda3\lib\site-packages\py4j\java_gateway.py:1321, in JavaMember.__call__(self, *args)
  1315 command = proto.CALL_COMMAND_NAME +\
  1316     self.command_header +\
  1317     args_command +\
  1318     proto.END_COMMAND_PART
  1320 answer = self.gateway_client.send_command(command)
-> 1321 return_value = get_return_value(
  1322     answer, self.gateway_client, self.target_id, self.name)
  1324 for temp_arg in temp_args:
  1325     temp_arg._detach()

文件 ~\anaconda3\lib\site-packages\pyspark\sql\utils.py:190, in capture_sql_exception.<locals>.deco(*a, **kw)
   188 def deco(*a: Any, **kw: Any) -> Any:
   189     try:
-> 190         return f(*a, **kw)
   191     except Py4JJavaError as e:
   192         converted = convert_exception(e.java_exception)

文件 ~\anaconda3\lib\site-packages\py4j\protocol.py:326, in get_return_value(answer, gateway_client, target_id, name)
   324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client)
   325 if answer[1] == REFERENCE_TYPE:
-> 326     raise Py4JJavaError(
   327         "调用 {0}.{1} 时发生错误。\n".
   328         format(target_id, name), value)
   329 else:
   330     raise Py4JError(
   331         "调用 {0}.{1} 时发生错误。堆栈跟踪:\n{3}\n".
   332         format(target_id, name, value))

Py4JJavaError: 调用 o46.showString 时发生错误。
: org.apache.spark.SparkException: 任务因阶段失败而中止:阶段0.0中的任务0失败1次,最近一次失败:丢失阶段0.0中的任务0.0(TID 0)(192.168.0.19 executor driver):org.apache.spark.SparkException: Python worker无法连接回驱动。
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:189)
   at org.apache.spark.api.python.PythonWorkerFactory.create(PythonWorkerFactory.scala:109)
   at org.apache.spark.SparkEnv.createPythonWorker(SparkEnv.scala:124)
   at org.apache.spark.api.python.BasePythonRunner.compute(PythonRunner.scala:164)
   at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
   at org.apache.spark.scheduler.Task.run(Task.scala:136)
   at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
   at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
   at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
   at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
   at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
   at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: java.net.SocketTimeoutException: 连接超时
   at java.base/java.net.PlainSocketImpl.waitForNewConnection(Native Method)
   at java.base/java.net.PlainSocketImpl.socketAccept(PlainSocketImpl.java:163)
   at java.base/java.net.AbstractPlainSocketImpl.accept(AbstractPlainSocketImpl.java:458)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:565)
   at java.base/java.net.ServerSocket.accept(ServerSocket.java:533)
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:176)
   ... 29 more

驱动堆栈跟踪:
   at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2672)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2608)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2607)
   at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
   at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
   at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
   at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2607)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1182)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1182)
   at scala.Option.foreach(Option.scala:407)
   at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1182)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2860)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2802)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2791)
   at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
   at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:952)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2228)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2249)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2268)
   at org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:506)
   at org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:459)
   at org.apache.spark.sql.execution.CollectLimitExec.executeCollect(limit.scala:48)
   at org.apache.spark.sql.Dataset.collectFromPlan(Dataset.scala:3868)
   at org.apache.spark.sql.Dataset.$anonfun$head$1(Dataset.scala:2863)
   at org.apache.spark.sql.Dataset.$anonfun$withAction$2(Dataset.scala:3858)
   at org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:510)
   at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:3856)
   at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:109)
   at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:169)
   at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:95)
   at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779)
   at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64)
   at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3856)
   at org.apache.spark.sql.Dataset.head(Dataset.scala:2863)
   at org.apache.spark.sql.Dataset.take(Dataset.scala:3084)
   at org.apache.spark.sql.Dataset.getRows(Dataset.scala:288)
   at org.apache.spark.sql.Dataset.showString(Dataset.scala:327)
   at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
   at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
   at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
   at java.base/java.lang.reflect.Method.invoke(Method.java:566)
   at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
   at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
   at py4j.Gateway.invoke(Gateway.java:282)
   at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
   at py4j.commands.CallCommand.execute(CallCommand.java:79)
   at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
   at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
   at java.base/java.lang.Thread.run(Thread.java:829)
Caused by: org.apache.spark.SparkException: Python worker无法连接回驱动。
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:189)
   at org.apache.spark.api.python.PythonWorkerFactory.create(PythonWorkerFactory.scala:109)
   at org.apache.spark.SparkEnv.createPythonWorker(SparkEnv.scala:124)
   at org.apache.spark.api.python.BasePythonRunner.compute(PythonRunner.scala:164)
   at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
   at org.apache.spark.scheduler.Task.run(Task.scala:136)
   at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
   at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
   at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
   at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
   at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
   ... 1 more
Caused by: java.net.SocketTimeoutException: 连接超时
   at java.base/java.net.PlainSocketImpl.waitForNewConnection(Native Method)
   at java.base/java.net.PlainSocketImpl.socketAccept(PlainSocketImpl.java:163)
   at java.base/java.net.AbstractPlainSocketImpl.accept(AbstractPlainSocketImpl.java:458)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:565)
   at java.base/java.net.ServerSocket.accept(ServerSocket.java:533)
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:176)
   ... 29 more

报错解决方法

报错核心是Python Worker无法与Spark驱动建立连接,属于PySpark环境配置问题,可尝试以下方案:

  • 指定Python解释器路径:在创建SparkSession时显式配置PySpark使用的Python路径(避免环境冲突):
spark = SparkSession.builder.appName("SparkTrial")\
    .config("spark.some.config.option", "some-value")\
    .config("spark.pyspark.python", "C:\\Users\\你的用户名\\anaconda3\\python.exe")\
    .config("spark.pyspark.driver.python", "C:\\Users\\你的用户名\\anaconda3\\python.exe")\
    .getOrCreate()
  • 调整连接超时时间:延长Worker与驱动的连接超时阈值:
spark = SparkSession.builder.appName("SparkTrial")\
    .config("spark.some.config.option", "some-value")\
    .config("spark.python.worker.connectionTimeout", "60s")\
    .getOrCreate()
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 11:20:27