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连接中断,常见触发因素包括:
- Spark Driver进程意外崩溃
- Py4J版本不兼容(Spark自带的Py4J版本与环境中的Python依赖冲突)
- 内存不足导致Driver或Executor被系统终止
- 防火墙/杀毒软件拦截了本地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
相关产品推荐
相关产品推荐

