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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 10:00:12