PySpark获取AWS S3子目录报错:Wrong FS问题排查求助
解决PySpark遍历S3子文件夹时的"Wrong FS"错误
问题原因
报错Wrong FS: s3a://prices/xml_inputs/stores, expected: file:///的核心原因有两个:
- 原代码直接通过
FileSystem.get()获取文件系统,默认绑定了本地文件系统(file:///),未正确关联S3A协议的配置 - Spark环境缺少
hadoop-aws依赖包,导致无法识别s3a协议对应的文件系统实现
修复方案
1. 补充依赖包
在SparkSession构建时添加对应版本的hadoop-aws包(版本需与环境中的Hadoop版本匹配,例如Hadoop 3.3.x对应3.3.4版本)。
2. 正确获取S3A文件系统
通过目标S3路径对应的Path对象来获取专属的FileSystem,确保使用的是S3A协议的文件系统实现。
修改后的完整代码
from pyspark.sql import SparkSession import os from dotenv import load_dotenv # 加载环境变量 load_dotenv() AWS_ACCESS_KEY_ID = os.getenv("AWS_ACCESS_KEY_ID") AWS_SECRET_ACCESS_KEY = os.getenv("AWS_SECRET_ACCESS_KEY") # 构建SparkSession并添加必要依赖 spark = SparkSession.builder \ .appName("price-xml-s3") \ .config("spark.jars.packages", "com.databricks:spark-xml_2.12:0.12.0,org.apache.hadoop:hadoop-aws:3.3.4") \ .config("spark.sql.execution.arrow.enabled", "true") \ .getOrCreate() # 配置S3相关参数 hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration() hadoop_conf.set("com.amazonaws.services.s3.enableV4", "true") hadoop_conf.set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") hadoop_conf.set("fs.s3a.aws.credentials.provider", "com.amazonaws.auth.InstanceProfileCredentialsProvider,com.amazonaws.auth.DefaultAWSCredentialsProviderChain") hadoop_conf.set("fs.AbstractFileSystem.s3a.impl", "org.apache.hadoop.fs.s3a.S3A") hadoop_conf.set("fs.s3a.access.key", AWS_ACCESS_KEY_ID) hadoop_conf.set("fs.s3a.secret.key", AWS_SECRET_ACCESS_KEY) hadoop_conf.set("fs.s3a.endpoint", "s3.us-east-1.amazonaws.com") # 定义目标S3路径 s3_path = "s3a://prices/xml_inputs/stores/" # 获取S3A文件系统实例 jvm = spark._jvm path = jvm.org.apache.hadoop.fs.Path(s3_path) fs = path.getFileSystem(hadoop_conf) # 遍历并筛选子文件夹 subfolders = [path.toString() for path in fs.listStatus(path) if path.isDirectory()] # 输出子文件夹(可选:仅提取文件夹名称) for subfolder in subfolders: folder_name = subfolder.rstrip('/').split('/')[-1] print(folder_name) # 关闭SparkSession spark.stop()
额外说明
- 若使用EMR等托管Spark环境,可移除
fs.s3a.access.key和fs.s3a.secret.key配置,依赖实例角色的权限即可 - 确保
hadoop-aws版本与Spark内置的Hadoop版本匹配,避免出现兼容性问题
内容的提问来源于stack exchange,提问作者Galat
相关产品推荐
相关产品推荐

