Windows本地Jupyter Notebook用PySpark读取S3 Parquet文件报错求助
问题描述
Windows本地已安装PySpark和Java,在Jupyter Notebook中尝试读取S3存储上的Parquet文件,执行代码如下:
import pyspark from pyspark.sql import SparkSession spark = SparkSession.builder \ .master('local') \ .appName('myAppName') \ .config('spark.executor.memory', '5gb') \ .config("spark.cores.max", "6") \ .getOrCreate() df = spark.read.parquet("s3://path/to/parquet/file.parquet")
运行后抛出错误:
Py4JJavaError: An error occurred while calling o30.parquet. : org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "s3" at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:3443) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3466) at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540) at org.apache.hadoop.fs.Path.getFileSystem(Path.java:365) at org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$checkAndGlobPathIfNecessary$1(DataSource.scala:752) at scala.collection.immutable.List.map(List.scala:293) at org.apache.spark.sql.execution.datasources.DataSource$.checkAndGlobPathIfNecessary(DataSource.scala:750) at org.apache.spark.sql.execution.datasources.DataSource.checkAndGlobPathIfNecessary(DataSource.scala:579) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:408) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:228) at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:210) at scala.Option.getOrElse(Option.scala:189) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:210) at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:562) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source) at java.lang.reflect.Method.invoke(Unknown Source) 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(Unknown Source)
解决方案
该错误源于PySpark缺少S3文件系统的支持依赖,同时未正确配置AWS访问凭证,按以下步骤处理:
1. 添加S3相关依赖包
创建SparkSession时,需指定匹配Spark/Hadoop版本的Hadoop AWS依赖包,示例如下:
spark = SparkSession.builder \ .master('local') \ .appName('myAppName') \ .config('spark.executor.memory', '5gb') \ .config("spark.cores.max", "6") \ # 替换为与你的Spark内置Hadoop版本匹配的依赖版本 .config("spark.jars.packages", "org.apache.hadoop:hadoop-aws:3.3.4,com.amazonaws:aws-java-sdk-bundle:1.12.262") \ # 指定S3文件系统实现类 .config("spark.hadoop.fs.s3.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") \ .getOrCreate()
可通过执行
spark-submit --version查看当前Spark依赖的Hadoop版本,确保hadoop-aws版本与之匹配。
2. 配置AWS访问凭证
提供两种配置方式:
- 代码内配置:创建SparkSession后添加凭证信息
spark._jsc.hadoopConfiguration().set("fs.s3a.access.key", "你的AWS_ACCESS_KEY_ID") spark._jsc.hadoopConfiguration().set("fs.s3a.secret.key", "你的AWS_SECRET_ACCESS_KEY") - 环境变量配置:在Windows命令提示符中设置环境变量,再启动Jupyter Notebook
set AWS_ACCESS_KEY_ID=你的AWS_ACCESS_KEY_ID set AWS_SECRET_ACCESS_KEY=你的AWS_SECRET_ACCESS_KEY
3. 修改S3路径前缀
将原代码中的s3://替换为Hadoop支持的s3a://协议前缀:
df = spark.read.parquet("s3a://path/to/parquet/file.parquet")
4. 手动安装依赖(可选)
若自动下载依赖失败,可手动下载对应版本的hadoop-aws.jar和aws-java-sdk-bundle.jar,放入Spark安装目录下的jars文件夹,重启Jupyter Notebook即可。
内容的提问来源于stack exchange,提问作者daniel____
相关产品推荐
相关产品推荐

