Spark3.4.0+Python3.11在Hadoop主节点执行spark.sql遇Py4JException错误
环境信息
- Apache Spark 3.4.0
- Python 3.11
- 运行环境:Hadoop主节点的Jupyter Notebook
问题描述
在主节点通过Jupyter Notebook运行PySpark连接Hive时,执行spark.sql("SHOW TABLES").show()触发Py4JError,提示py4j.Py4JException: Method sql([class java.lang.String, class [Ljava.lang.Object;]) does not exist,其他虚拟机可正常连接Hive,但必须使用主节点。
代码示例
from pyspark.sql import SparkSession from pyspark.sql import Row spark = SparkSession.builder \ .appName("YourAppName") \ .config("spark.sql.hive.hiveserver2.jdbc.url", "jdbc:hive2://{IP}:{PORT}/{SCHEMA};user={USER};password={PWD}") \ .config("spark.master", "spark://{IP}:{PORT}") \ .enableHiveSupport() \ .getOrCreate() spark.sql("SHOW TABLES").show()
错误栈信息
Py4JError Traceback (most recent call last) Cell In[3], line 1 ----> 1 spark.sql("SHOW TABLES").show() File ~/anaconda3/lib/python3.11/site-packages/pyspark/sql/session.py:1631, in SparkSession.sql(self, sqlQuery, args, **kwargs) 1627 assert self._jvm is not None 1628 litArgs = self._jvm.PythonUtils.toArray( 1629 [_to_java_column(lit(v)) for v in (args or [])] 1630 ) -> 1631 return DataFrame(self._jsparkSession.sql(sqlQuery, litArgs), self) 1632 finally: 1633 if len(kwargs) > 0: File ~/anaconda3/lib/python3.11/site-packages/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 ~/anaconda3/lib/python3.11/site-packages/pyspark/errors/exceptions/captured.py:179, in capture_sql_exception.<locals>.deco(*a, **kw) 177 def deco(*a: Any, **kw: Any) -> Any: 178 try: --> 179 return f(*a, **kw) 180 except Py4JJavaError as e: 181 converted = convert_exception(e.java_exception) File ~/anaconda3/lib/python3.11/site-packages/py4j/protocol.py:330, in get_return_value(answer, gateway_client, target_id, name) 326 raise Py4JJavaError( 327 "An error occurred while calling {0}{1}{2}.\n". 328 format(target_id, ".", name), value) 329 else: --> 330 raise Py4JError( 331 "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n". 332 format(target_id, ".", name, value)) 333 else: 334 raise Py4JError( 335 "An error occurred while calling {0}{1}{2}". 336 format(target_id, ".", name)) Py4JError: An error occurred while calling o40.sql. Trace: py4j.Py4JException: Method sql([class java.lang.String, class [Ljava.lang.Object;]) does not exist at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:321) at py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:329) at py4j.Gateway.invoke(Gateway.java:274) 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.lang.Thread.run(Thread.java:750)
核心原因
该错误源于Python端PySpark API与Java端Spark内核的版本/参数签名不匹配,或是Jupyter使用的Python环境与Spark绑定的Python版本不一致,也可能是Hive配置加载异常导致的方法调用失败。
具体修复步骤
对齐Python版本
Spark 3.4.0官方仅支持Python 3.7-3.10,Python 3.11存在兼容性风险,建议降级到Python 3.10。执行以下命令确认版本匹配:python --version spark-submit --version修正SparkSession配置
spark.sql.hive.hiveserver2.jdbc.url是JDBC连接HiveServer2的参数,而enableHiveSupport()是让Spark读取本地Hive元数据配置,无需该参数。修改后的代码:from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("YourAppName") \ .config("spark.master", "spark://{IP}:{PORT}") \ .enableHiveSupport() \ .getOrCreate() spark.sql("SHOW TABLES").show()同时确保Hive的
hive-site.xml已放入Spark的conf目录,或启动Jupyter时通过--files参数加载该配置。配置主节点环境变量
在主节点设置正确的环境变量后重启Jupyter:export SPARK_HOME=/path/to/spark-3.4.0 export PYSPARK_PYTHON=/path/to/python3.10 export PYSPARK_DRIVER_PYTHON=jupyter export PYSPARK_DRIVER_PYTHON_OPTS="notebook"验证Spark内核兼容性
在主节点执行spark-shell,运行sql("SHOW TABLES").show(),如果同样报错则是Spark内核问题,否则排查Jupyter环境。清理旧会话
在Jupyter中先停止旧的SparkSession再重新创建:spark.stop() # 重新创建SparkSession
内容的提问来源于stack exchange,提问作者Dishan

