使用pytest测试PySpark保存DataFrame为表时遇Py4JJavaError求助
PySpark测试中saveAsTable抛出Py4JJavaError问题排查
我是pytest新手,正在测试一段PySpark代码:读取parquet文件,将double类型转换为decimal类型后追加到表中。但执行df.write.mode("append").saveAsTable("tablename")时测试用例失败,抛出Py4JJavaError。以下是我的pytest代码、主类代码及完整错误堆栈信息,恳请帮忙排查问题:
pytest代码
@pytest.fixture(scope="module") def sample_df(spark): data = [(1,2.0, 12,'Test')] schema = StructType( [ StructField("foo1", IntegerType(), True), StructField("foo2", DoubleType(), True), StructField("foo3", IntegerType(), True), StructField("Name", StringType(), True) ] ) sample_df = spark.createDataFrame(data, schema) sample_df.createOrReplaceTempView("dummy_table") return sample_df @pytest.mark.parametrize('dbutils', [run_params], indirect=['dbutils']) def test_loadDataIntoForecastTable( spark_mock, dbutils,sample_df): spark_mock.read.parquet.return_value = sample_df loadDataIntoForecastTable(spark_mock,"test","test","test") spark_mock.sql.assert_called()
主类代码
def loadDataIntoForecastTable(spark,bar1,bar2, s3_file_name): df = spark.read.parquet(s3_file_name) for dtype in df.dtypes: if(dtype[1] == 'double'): column = dtype[0] df = df.withColumn(column,df[column].cast("decimal(38,9)")) #delete existing data from the table df.select("foo1").distinct().createOrReplaceTempView("tempTableDeleteExisting") delete_query = """DELETE FROM {} fha WHERE ID in (SELECT foo1 from tempTableDeleteExisting)""" delete_sql = delete_query.format(bar2) print("the delete sql is : ",delete_sql) spark.sql(delete_sql) #update table df.write.mode("append").option("mergeSchema", "true").saveAsTable(bar2)
错误堆栈
> answer = 'xro160', gateway_client = <py4j.clientserver.JavaClient > object at 0x00000279B6566260>, target_id = 'o159', name = > 'saveAsTable' > > def get_return_value(answer, gateway_client, target_id=None, name=None): > """Converts an answer received from the Java gateway into a Python object. > > For example, string representation of integers are converted to Python > integer, string representation of objects are converted to JavaObject > instances, etc. > > :param answer: the string returned by the Java gateway > :param gateway_client: the gateway client used to communicate with the Java > Gateway. Only necessary if the answer is a reference (e.g., object, > list, map) > :param target_id: the name of the object from which the answer comes from > (e.g., *object1* in `object1.hello()`). Optional. > :param name: the name of the member from which the answer comes from > (e.g., *hello* in `object1.hello()`). Optional. > """ > if is_error(answer)[0]: > if len(answer) > 1: > type = answer[1] > value = OUTPUT_CONVERTER[type](answer[2:], gateway_client) > if answer[1] == REFERENCE_TYPE: > > raise Py4JJavaError( > "An error occurred while calling {0}{1}{2}.\n". > format(target_id, ".", name), value) E py4j.protocol.Py4JJavaError: An error occurred while calling > o159.saveAsTable. E : > java.lang.UnsatisfiedLinkError: > org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z > E at > org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Native Method) > E at > org.apache.hadoop.io.nativeio.NativeIO$Windows.access(NativeIO.java:793) > E at > org.apache.hadoop.fs.FileUtil.canRead(FileUtil.java:1249) E > at org.apache.hadoop.fs.FileUtil.list(FileUtil.java:1454) E > at > org.apache.hadoop.fs.RawLocalFileSystem.listStatus(RawLocalFileSystem.java:601) > E at > org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972) E > at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014) E > at > org.apache.hadoop.fs.ChecksumFileSystem.listStatus(ChecksumFileSystem.java:761) > E at > org.apache.spark.sql.catalyst.catalog.SessionCatalog.validateTableLocation(SessionCatalog.scala:413) > E at > org.apache.spark.sql.execution.command.CreateDataSourceTableAsSelectCommand.run(createDataSourceTables.scala:176) > E at > org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:75) > E at > org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:73) > E at > org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:84) > E at > org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:107) > E at > org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:125) > E at > org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:201) > E at > org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:108) > E at > org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:900) E > at > org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66) > E at > org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:107) > E at > org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:98) > E at > org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:461) > E at > org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(origin.scala:76) > E at > org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:461) > E at > org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:32) > E at > org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) > E at > org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) > E at > org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:32) > E at > org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:32) > E at > org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:437) > E at > org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:98) > E at > org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:85) > E at > org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:83) > E at > org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:142) > E at > org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:859) > E at > org.apache.spark.sql.DataFrameWriter.createTable(DataFrameWriter.scala:700) > E at > org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:678) > E at > org.apache.spark.sql.DataFrameWriter.saveAsTable(DataFrameWriter.scala:571) > E at > java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native > Method) E at > java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) > E at > java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) > E at > java.base/java.lang.reflect.Method.invoke(Method.java:566) E > at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) E > at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) > E at py4j.Gateway.invoke(Gateway.java:282) E > at > py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) E > at py4j.commands.CallCommand.execute(CallCommand.java:79) E > at > py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) > E at > py4j.ClientServerConnection.run(ClientServerConnection.java:106) E > at java.base/java.lang.Thread.run(Thread.java:834) > > C:\JPMC\DEV\TMP\ds\tools\python3.10\latest\lib\site-packages\py4j\protocol.py:326: > Py4JJavaError
问题排查与解决
核心原因
错误堆栈中的java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z是关键:这是Windows环境下Hadoop原生库不兼容导致的——你的Hadoop版本与当前Windows系统的JDK版本不匹配,或者缺少对应的Windows原生库(hadoop.dll、winutils.exe)。
同时测试代码和主代码还有两个潜在问题:
- 测试中仅mock了
spark.read.parquet,但未mocksaveAsTable和spark.sql的真实执行逻辑,导致测试触发了Spark的实际文件系统操作,暴露了Hadoop原生库问题。 - 主代码的DELETE语句存在字段不匹配:DataFrame中是
foo1,但SQL里用了ID,如果目标表没有ID字段,运行时会触发额外报错。
解决步骤
1. 修复Windows下Hadoop原生库问题
- 下载与你的Hadoop版本匹配的Windows原生工具包(包含hadoop.dll、winutils.exe)。
- 将这些文件放到Hadoop安装目录的
bin文件夹下。 - 配置系统环境变量
HADOOP_HOME指向你的Hadoop安装目录,同时将%HADOOP_HOME%\bin添加到PATH中。 - 重启Python运行环境,确保环境变量生效。
2. 完善测试用例,mock所有Spark操作
单元测试应完全mock掉IO相关操作,避免触发真实执行:
@pytest.mark.parametrize('dbutils', [run_params], indirect=['dbutils']) def test_loadDataIntoForecastTable( spark_mock, dbutils,sample_df): spark_mock.read.parquet.return_value = sample_df # mock spark.sql方法 spark_mock.sql.return_value = None # 链式mock DataFrame的write操作 mock_df_write = spark_mock.read.parquet.return_value.write mock_df_write.mode.return_value = mock_df_write mock_df_write.option.return_value = mock_df_write mock_df_write.saveAsTable.return_value = None loadDataIntoForecastTable(spark_mock,"test","test","test") spark_mock.sql.assert_called() mock_df_write.saveAsTable.assert_called_with("test")
3. 修复主代码中的SQL字段错误
将DELETE语句中的ID改为foo1,确保与DataFrame字段一致:
delete_query = """DELETE FROM {} fha WHERE foo1 in (SELECT foo1 from tempTableDeleteExisting)"""
内容的提问来源于stack exchange,提问作者Ananya
相关产品推荐
相关产品推荐

