Windows11下PySpark写入CSV触发Py4JJavaError问题求助
PySpark写入CSV时触发Py4JJavaError(根源为UnsatisfiedLinkError)
问题描述
运行PySpark脚本读取CSV并完成过滤、聚合操作后,执行write.csv时触发Py4JJavaError,此前数据读取、转换步骤均正常运行。已尝试将Spark目录迁移至用户目录并更新环境变量,问题未解决。
运行环境
Windows 11,初始环境变量配置:
setx JAVA_HOME “C:\Program Files\Java\jdk1.8.0_202” setx SPARK_HOME “C:\spark-3.3.2-bin-hadoop3” setx HADOOP_HOME “C:\spark-3.3.2-bin-hadoop3” setx PATH “%PATH%;C:\spark-3.3.2-bin-hadoop3\bin”
PySpark代码
import findspark findspark.init() findspark.find() from pyspark.sql import SparkSession # 创建SparkSession spark = SparkSession.builder.appName("large_dataset_example").getOrCreate() # 从CSV读取DataFrame df = spark.read.csv("data.csv", header=True, inferSchema=True) # 展示前5行 df.show(5) # 过滤数据 filtered_df = df.filter((df.age >= 30) & (df.gender == "Male")) # 展示过滤结果 filtered_df.show() # 按性别分组计算平均年龄 avg_age_df = filtered_df.groupBy("gender").agg({"age": "avg"}) # 展示平均年龄结果 avg_age_df.show() ##################################################################### # 将结果写入新CSV filtered_df.write.csv("filtered_data.csv", header=True) avg_age_df.write.csv("average_age.csv", header=True)
错误堆栈
初始错误(根源)
--------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call last) Cell In[9], line 2 1 # Write the filtered and aggregated data to a new CSV file ----> 2 filtered_df.write.csv("filtered_data.csv", header=True) 3 avg_age_df.write.csv("average_age.csv", header=True) File C:\spark-3.3.2-bin-hadoop3\python\pyspark\sql\readwriter.py:1240, in DataFrameWriter.csv(self, path, mode, compression, sep, quote, escape, header, nullValue, escapeQuotes, quoteAll, dateFormat, timestampFormat, ignoreLeadingWhiteSpace, ignoreTrailingWhiteSpace, charToEscapeQuoteEscaping, encoding, emptyValue, lineSep) 1221 self.mode(mode) 1222 self._set_opts( 1223 compression=compression, 1224 sep=sep, (...) 1238 lineSep=lineSep, 1239 ) -> 1240 self._jwrite.csv(path) File C:\spark-3.3.2-bin-hadoop3\python\lib\py4j-0.10.9.5-src.zip\py4j\java_gateway.py:1321, in JavaMember.__call__(self, *args) 1315 command = proto.CALL_COMMAND_NAME +\ 1316 self.command_header +\ 1317 args_command +\ 1318 proto.END_COMMAND_PART 1320 answer = self.gateway_client.send_command(command) -> 1321 return_value = get_return_value( 1322 answer, self.gateway_client, self.target_id, self.name) 1324 for temp_arg in temp_args: 1325 temp_arg._detach() File C:\spark-3.3.2-bin-hadoop3\python\pyspark\sql\utils.py:190, in capture_sql_exception.<locals>.deco(*a, **kw) 188 def deco(*a: Any, **kw: Any) -> Any: 189 try: --> 190 return f(*a, **kw) 191 except Py4JJavaError as e: 192 converted = convert_exception(e.java_exception) File C:\spark-3.3.2-bin-hadoop3\python\lib\py4j-0.10.9.5-src.zip\py4j\protocol.py:326, in get_return_value(answer, gateway_client, target_id, name) 324 value = OUTPUT_CONVERTER[type](answer[2:], gateway_client) 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: 330 raise Py4JError( 331 "An error occurred while calling {0}{1}{2}. Trace:\n{3}\n". 332 format(target_id, ".", name, value)) Py4JJavaError: An error occurred while calling o68.csv. : org.apache.spark.SparkException: Job aborted. at org.apache.spark.sql.errors.QueryExecutionErrors$.jobAbortedError(QueryExecutionErrors.scala:651) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.write(FileFormatWriter.scala:288) at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:186) at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult$lzycompute(commands.scala:113) at org.apache.spark.sql.execution.command.DataWritingCommandExec.sideEffectResult(commands.scala:111) at org.apache.spark.sql.execution.command.DataWritingCommandExec.executeCollect(commands.scala:125) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.$anonfun$applyOrElse$1(QueryExecution.scala:98) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:109) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:169) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:95) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:98) at org.apache.spark.sql.execution.QueryExecution$$anonfun$eagerlyExecuteCommands$1.applyOrElse(QueryExecution.scala:94) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDownWithPruning$1(TreeNode.scala:584) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:176) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDownWithPruning(TreeNode.scala:584) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning(AnalysisHelper.scala:267) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDownWithPruning$(AnalysisHelper.scala:263) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDownWithPruning(LogicalPlan.scala:30) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:560) at org.apache.spark.sql.execution.QueryExecution.eagerlyExecuteCommands(QueryExecution.scala:94) at org.apache.spark.sql.execution.QueryExecution.commandExecuted$lzycompute(QueryExecution.scala:81) at org.apache.spark.sql.execution.QueryExecution.commandExecuted(QueryExecution.scala:79) at org.apache.spark.sql.execution.QueryExecution.assertCommandExecuted(QueryExecution.scala:116) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:860) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:390) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:363) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:239) at org.apache.spark.sql.DataFrameWriter.csv(DataFrameWriter.scala:851) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) 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:748) Caused by: java.lang.UnsatisfiedLinkError: org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Ljava/lang/String;I)Z at org.apache.hadoop.io.nativeio.NativeIO$Windows.access0(Native Method) at org.apache.hadoop.io.nativeio.NativeIO$Windows.access(NativeIO.java:793) at org.apache.hadoop.fs.FileUtil.canRead(FileUtil.java:1218) at org.apache.hadoop.fs.FileUtil.list(FileUtil.java:1423) at org.apache.hadoop.fs.RawLocalFileSystem.listStatus(RawLocalFileSystem.java:601) at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972) at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014) at org.apache.hadoop.fs.ChecksumFileSystem.listStatus(ChecksumFileSystem.java:761) at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1972) at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:2014) at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.getAllCommittedTaskPaths(FileOutputCommitter.java:334) at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJobInternal(FileOutputCommitter.java:404) at org.apache.hadoop.mapreduce.lib.output.FileOutputCommitter.commitJob(FileOutputCommitter.java:377) at org.apache.spark.internal.io.HadoopMapReduceCommitProtocol.commitJob(HadoopMapReduceCommitProtocol.scala:192) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.$anonfun$write$26(FileFormatWriter.scala:277) at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.util.Utils$.timeTakenMs(Utils.scala:642) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.write(FileFormatWriter.scala:277) ... 42 more
调整目录后的错误
Output exceeds the size limit. Open the full output data in a text editor --------------------------------------------------------------------------- Py4JJavaError Traceback (most recent call last) Cell In[17], line 2 1 # Write the filtered and aggregated data to a new CSV file ----> 2 filtered_df.write.csv("filtered_data.csv", header=True) 3 avg_age_df.write.csv("average_age.csv", header=True) File ~\spark-3.3.2-bin-hadoop3\python\pyspark\sql\readwriter.py:1240, in DataFrameWriter.csv(self, path, mode, compression, sep, quote, escape, header, nullValue, escapeQuotes, quoteAll, dateFormat, timestampFormat, ignoreLeadingWhiteSpace, ignoreTrailingWhiteSpace, charToEscapeQuoteEscaping, encoding, emptyValue, lineSep) 1221 self.mode(mode) 1222 self._set_opts( 1223 compression=compression, 1224 sep=sep, (...) 1238 lineSep=lineSep, 1239 ) -> 1240 self._jwrite.csv(path) File ~\spark-3.3.2-bin-hadoop3\python\lib\py4j-0.10.9.5-src.zip\py4j\java_gateway.py:1321, in JavaMember.__call__(self, *args) 1315 command = proto.CALL_COMMAND_NAME +\ 1316 self.command_header +\ 1317 args_command +\ 1318 proto.END_COMMAND_PART 1320 answer = self.gateway_client.send_command(command) -> 1321 return_value = get_return_value( 1322 answer, self.gateway_client, self.target_id, self.name) ... at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) at org.apache.spark.util.Utils$.timeTakenMs(Utils.scala:642) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.write(FileFormatWriter.scala:277) ... 42 more
测试数据集(10行)
id,name,age,gender 1,John,30,Male 2,Jane,25,Female 3,Mark,40,Male 4,Sara,35,Female 5,David,28,Male 6,Emily,32,Female 7,Steven,45,Male 8,Amy,27,Female 9,Chris,50,Male 10,Lisa,42,Female
解决方法
该错误是Windows环境下Hadoop原生IO依赖缺失或权限问题导致,以下是几种可行修复方案:
方案1:补充Hadoop Windows原生库
- 下载与Spark配套Hadoop版本(3.x)对应的Windows原生库文件(
winutils.exe、hadoop.dll) - 在Spark根目录下创建
bin文件夹(若无),将下载的文件放入该目录 - 确认
HADOOP_HOME指向Spark目录,且%HADOOP_HOME%\bin已添加至系统PATH
方案2:禁用Hadoop原生IO检查
在初始化SparkSession时添加配置,强制使用非原生IO实现:
spark = SparkSession.builder \ .appName("large_dataset_example") \ .config("spark.hadoop.fs.file.impl", "org.apache.hadoop.fs.LocalFileSystem") \ .config("spark.hadoop.hadoop.security.authorization", "false") \ .getOrCreate()
方案3:以管理员权限运行脚本
Windows系统权限限制可能触发该错误,尝试以管理员身份启动Python/Notebook环境后再执行脚本。
方案4:调整输出路径权限
避免将结果写入系统保护目录(如C盘根目录),选择用户目录下的可读写文件夹作为输出路径,确保当前用户对目标目录有完整读写权限。
内容的提问来源于stack exchange,提问作者user8314628
相关产品推荐
相关产品推荐

