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

使用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)。

同时测试代码和主代码还有两个潜在问题:

  1. 测试中仅mock了spark.read.parquet,但未mocksaveAsTable和spark.sql的真实执行逻辑,导致测试触发了Spark的实际文件系统操作,暴露了Hadoop原生库问题。
  2. 主代码的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:55:55