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

Spark DataFrame重分区后出现ConnectionRefusedError问题求助

Spark重分区后执行show()出现Py4J连接错误

我的Spark DataFrame包含12个分区(2022年纽约黄色出租车行程记录数据),为平衡分区大小,执行了以下重分区操作:

taxi_df = taxi_df.repartition(10)

但重分区后运行taxi_df.show()时抛出如下异常:

Py4JJavaError                             Traceback (most recent call last)
File ~\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\errors\exceptions\captured.py:179, in capture_sql_exception.<locals>.deco(*a, **kw)
   178 try:
--> 179     return f(*a, **kw)
   180 except Py4JJavaError as e:

File ~\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\protocol.py:326, in get_return_value(answer, gateway_client, target_id, name)
   325 if answer[1] == REFERENCE_TYPE:
--> 326     raise Py4JJavaError(
   327         "An error occurred while calling {0}{1}{2}.\n".
   328         format(target_id, ".", name), value)
   329 else:

<class 'str'>: (<class 'ConnectionRefusedError'>, ConnectionRefusedError(10061, 'No connection could be made because the target machine actively refused it', None, 10061, None))

During handling of the above exception, another exception occurred:

Py4JError                                 Traceback (most recent call last)
Cell In[10], line 1
----> 1 taxi_df.show()

File ~\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\sql\dataframe.py:945, in DataFrame.show(self, n, truncate, vertical)
   885 def show(self, n: int = 20, truncate: Union[bool, int] = True, vertical: bool = False) -> None:
   886     """Prints the first ``n`` rows to the console.
   887 
   888     .. versionadded:: 1.3.0
  (...)
   943     name | Bob
   944     """
--> 945     print(self._show_string(n, truncate, vertical))

File ~\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\sql\dataframe.py:963, in DataFrame._show_string(self, n, truncate, vertical)
   957     raise PySparkTypeError(
   958         error_class="NOT_BOOL",
   959         message_parameters={"arg_name": "vertical", "arg_type": type(vertical).__name__},
   960     )
   962 if isinstance(truncate, bool) and truncate:
--> 963     return self._jdf.showString(n, 20, vertical)
   964 else:
   965     try:

File ~\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\java_gateway.py:1322, in JavaMember.__call__(self, *args)
  1316 command = proto.CALL_COMMAND_NAME +\
  1317     self.command_header +\
  1318     args_command +\
  1319     proto.END_COMMAND_PART
  1321 answer = self.gateway_client.send_command(command)
--> 1322 return_value = get_return_value(
  1323     answer, self.gateway_client, self.target_id, self.name)
  1325 for temp_arg in temp_args:
  1326     if hasattr(temp_arg, "_detach"):

File ~\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\errors\exceptions\captured.py:181, in capture_sql_exception.<locals>.deco(*a, **kw)
   179     return f(*a, **kw)
   180 except Py4JJavaError as e:
--> 181     converted = convert_exception(e.java_exception)
   182     if not isinstance(converted, UnknownException):
   183         # Hide where the exception came from that shows a non-Pythonic
   184         # JVM exception message.
   185         raise converted from None

File ~\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\errors\exceptions\captured.py:143, in convert_exception(e)
   141 elif is_instance_of(gw, e, "java.lang.IllegalArgumentException"):
   142     return IllegalArgumentException(origin=e)
--> 143 elif is_instance_of(gw, e, "java.lang.ArithmeticException"):
   144     return ArithmeticException(origin=e)
   145 elif is_instance_of(gw, e, "java.lang.UnsupportedOperationException"):

File ~\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\java_gateway.py:464, in is_instance_of(gateway, java_object, java_class)
   460 else:
   461     raise Py4JError(
   462         "java_class must be a string, a JavaClass, or a JavaObject")
--> 464 return gateway.jvm.py4j.reflection.TypeUtil.isInstanceOf(
   465     param, java_object)

File ~\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\java_gateway.py:1664, in JavaPackage.__getattr__(self, name)
  1661     return JavaClass(
  1662         answer[proto.CLASS_FQN_START:], self._gateway_client)
  1663 else:
--> 1664     raise Py4JError("{0} does not exist in the JVM".format(new_fqn))

Py4JError: py4j.reflection does not exist in the JVM

此外,还有以下仅打印未抛出的错误信息:

----------------------------------------
Exception occurred during processing of request from ('127.0.0.1', 60537)
Traceback (most recent call last):
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socketserver.py", line 317, in _handle_request_noblock
    self.process_request(request, client_address)
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socketserver.py", line 348, in process_request
    self.finish_request(request, client_address)
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socketserver.py", line 361, in finish_request
    self.RequestHandlerClass(request, client_address, self)
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socketserver.py", line 755, in __init__
    self.handle()
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\accumulators.py", line 295, in handle
    poll(accum_updates)
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\accumulators.py", line 267, in poll
    if self.rfile in r and func():
                           ^^^^^^
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\accumulators.py", line 271, in accum_updates
    num_updates = read_int(self.rfile)
                  ^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\serializers.py", line 594, in read_int
    length = stream.read(4)
             ^^^^^^^^^^^^^^
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socket.py", line 706, in readinto
    return self._sock.recv_into(b)
           ^^^^^^^^^^^^^^^^^^^^^^^
ConnectionResetError: [WinError 10054] An existing connection was forcibly closed by the remote host
----------------------------------------
ERROR:root:Exception while sending command.
Traceback (most recent call last):
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\pyspark\errors\exceptions\captured.py", line 179, in deco
    return f(*a, **kw)
           ^^^^^^^^^^^
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\protocol.py", line 326, in get_return_value
    raise Py4JJavaError(
py4j.protocol.Py4JJavaError: <exception str() failed>

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\clientserver.py", line 511, in send_command
    answer = smart_decode(self.stream.readline()[:-1])
                          ^^^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\jatin\AppData\Local\Programs\Python\Python311\Lib\socket.py", line 706, in readinto
    return self._sock.recv_into(b)
           ^^^^^^^^^^^^^^^^^^^^^^^
ConnectionResetError: [WinError 10054] An existing connection was forcibly closed by the remote host

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\java_gateway.py", line 1038, in send_command
    response = connection.send_command(command)
               ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "C:\Users\jatin\Spark\spark-3.5.1-bin-hadoop3\python\lib\py4j-0.10.9.7-src.zip\py4j\clientserver.py", line 539, in send_command
    raise Py4JNetworkError(
py4j.protocol.Py4JNetworkError: Error while sending or receiving

问题分析与解决办法

核心原因

这些错误本质是Python与Spark JVM之间的Py4J连接中断,常见触发因素包括:

  1. Spark Driver进程意外崩溃
  2. Py4J版本不兼容(Spark自带的Py4J版本与环境中的Python依赖冲突)
  3. 内存不足导致Driver或Executor被系统终止
  4. 防火墙/杀毒软件拦截了本地Py4J通信端口

具体修复步骤

  • 重启Spark会话:直接关闭当前Notebook/终端,重新启动SparkContext/SparkSession,这是最快速的临时修复方案。
  • 检查Py4J版本兼容性:确保环境中没有单独安装Py4J,Spark自带的Py4J版本已经适配当前Spark版本,执行pip uninstall py4j移除独立安装的版本。
  • 调整内存配置:纽约出租车数据量较大,重分区操作可能消耗大量内存。在启动Spark时增加Driver内存:
    from pyspark.sql import SparkSession
    spark = SparkSession.builder \
        .appName("TaxiDataProcessing") \
        .config("spark.driver.memory", "8g") \
        .getOrCreate()
    
  • 检查系统防火墙/杀毒软件:临时关闭防火墙或添加Spark进程到信任列表,排除端口拦截问题。
  • 避免直接show()大分区数据:可以先采样查看数据,减少内存压力:
    taxi_df.sample(fraction=0.01).show()
    
  • 验证重分区操作:先执行taxi_df.rdd.getNumPartitions()确认分区数是否正确,再进行后续操作。

内容的提问来源于stack exchange,提问作者Jatin Rathour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 19:20:55