Spark Dataframe操作遇Py4JJavaError及连接拒绝问题求助
问题描述
执行聚合后的Spark DataFrame的count()操作,或转换为Pandas DataFrame时,出现如下错误:
File "C:\Spark\python\pyspark\sql\dataframe.py", line 804, in count return int(self._jdf.count()) File "C:\Users\ariel\AppData\Roaming\Python\Python39\site-packages\py4j\java_gateway.py", line 1321, in __call__ return_value = get_return_value( File "C:\Spark\python\pyspark\sql\utils.py", line 190, in deco return f(*a, **kw) File "C:\Users\ariel\AppData\Roaming\Python\Python39\site-packages\py4j\protocol.py", line 326, in get_return_value raise Py4JJavaError( py4j.protocol.Py4JJavaError: <unprintable Py4JJavaError object> During handling of the above exception, another exception occurred: response = connection.send_command(command) File "C:\Users\ariel\AppData\Roaming\Python\Python39\site-packages\py4j\clientserver.py", line 539, in send_command raise Py4JNetworkError( py4j.protocol.Py4JNetworkError: Error while sending or receiving ConnectionRefusedError: [WinError 10061] 由于目标计算机积极拒绝,无法建立连接
已尝试调整Spark默认内存至10GB、重装Spark和Py4J,问题仍存在,输入数据量小于10GB。
解决方案
- 检查Spark进程状态:出现连接拒绝通常意味着Spark Driver或Worker进程已崩溃。打开任务管理器查看Java进程是否存活,或访问Spark UI(默认端口4040)确认服务状态。若进程崩溃,查找Spark日志(本地模式下日志多位于临时目录或指定日志路径),定位OOM或分区数据过大等崩溃原因。
- 优化分区与聚合配置:
- 聚合前对数据重分区,避免单分区数据过载:
df = df.repartition(10)(分区数根据实际数据量调整) - 调整Shuffle分区数:
spark.conf.set("spark.sql.shuffle.partitions", 50)(默认200,减少不必要的Shuffle开销) - 启用广播连接(若涉及Join操作):
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10485760")(设置为10MB,小表自动广播)
- 聚合前对数据重分区,避免单分区数据过载:
- 调整JVM内存配置:除Executor内存外,需确保Driver有足够内存接收聚合结果。初始化SparkSession时显式配置:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .config("spark.driver.memory", "8g") \ .config("spark.executor.memory", "10g") \ .getOrCreate() - 验证版本兼容性:确认Python 3.9与所安装的Spark版本兼容(Spark 3.2及以上版本支持Python 3.9),版本不匹配可能导致进程异常退出。
- 释放系统资源:Windows本地模式下,关闭其他高内存占用程序,或调整系统虚拟内存大小,避免系统资源耗尽导致Spark进程崩溃。
内容的提问来源于stack exchange,提问作者miss_ariel
相关产品推荐
相关产品推荐

