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

PySpark写入S3后重命名part文件报错,EMR环境求解决方案

问题

在EMR环境下使用PySpark创建DataFrame并写入S3路径,生成part-*格式文件。尝试通过Hadoop文件工具库重命名文件时持续报错,且无法使用boto3。请问该操作在PySpark中是否可行?

使用的PySpark代码:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("test_rename").getOrCreate()
sc = spark.sparkContext

l = [['Column1', 'Column2', 'Column3'], ['Col1Value1', 'Col2Value1', 'Col3Value1'], ['Col1Value2', 'Col2Value2', 'Col3Value2']]
# 读取字符串数据为DataFrame
df = spark.createDataFrame(l[1:],l[0])

df.coalesce(1) \
    .write.format("com.databricks.spark.csv") \
    .option("header", "true") \
    .mode("overwrite") \
    .save("s3://<bucket>/<prefix>")

from py4j.java_gateway import java_import
java_import(spark._jvm, 'org.apache.hadoop.fs.Path')

fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
file = fs.globStatus(sc._jvm.Path('s3://<bucket>/<prefix>/part*'))[0].getPath().getName()
fs.rename(sc._jvm.Path('s3://<bucket>/<prefix>/'+ file), sc._jvm.Path('mydata.csv'))
fs.delete(sc._jvm.Path('s3://<bucket>/<prefix>'), True)

报错信息:

File "/mnt/tmp/spark-471166fb-d7c7-4839-a308-2e3f5c01c185/test_rename.py", line 20, in <module>
    file = fs.globStatus(sc._jvm.Path('s3://<bucket>/<prefix>/part*'))[0].getPath().getName()
  File "/usr/lib/spark/python/lib/py4j-0.10.9.5-src.zip/py4j/java_gateway.py", line 1322, in __call__
  File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 196, in deco
pyspark.sql.utils.IllegalArgumentException: Wrong FS: s3://<bucket>/<prefix>, expected: hdfs://<emr-ip>:8020
解决方案

该操作完全可行,报错核心原因是你获取的FileSystem实例默认是EMR的HDFS文件系统,无法处理S3路径。只需调整FileSystem的获取方式,就能正常操作S3上的文件。

问题根源

FileSystem.get()方法会根据Hadoop配置返回默认文件系统(这里是EMR集群的HDFS),而S3路径需要使用专门的S3文件系统实现(如S3AFileSystem),用HDFS实例直接操作S3路径就会触发"Wrong FS"错误。

修复后的代码

通过S3路径对应的Path实例获取正确的FileSystem,同时确保目标路径是完整的S3路径:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("test_rename").getOrCreate()
sc = spark.sparkContext

l = [['Column1', 'Column2', 'Column3'], ['Col1Value1', 'Col2Value1', 'Col3Value1'], ['Col1Value2', 'Col2Value2', 'Col3Value2']]
df = spark.createDataFrame(l[1:], l[0])

# 写入S3
df.coalesce(1) \
    .write.format("com.databricks.spark.csv") \
    .option("header", "true") \
    .mode("overwrite") \
    .save("s3://<bucket>/<prefix>")

from py4j.java_gateway import java_import
java_import(spark._jvm, 'org.apache.hadoop.fs.Path')

# 关键修复:通过S3路径获取对应FileSystem
conf = spark._jsc.hadoopConfiguration()
s3_part_path = spark._jvm.Path("s3://<bucket>/<prefix>/part*")
fs = s3_part_path.getFileSystem(conf)

# 查找part文件
status_list = fs.globStatus(s3_part_path)
if status_list:
    source_path = status_list[0].getPath()
    # 目标路径必须指定完整S3路径,避免写到本地/HDFS
    target_path = spark._jvm.Path("s3://<bucket>/<prefix>/mydata.csv")
    # 执行重命名
    if fs.rename(source_path, target_path):
        # 删除原文件夹(包含_SUCCESS等临时文件)
        fs.delete(spark._jvm.Path("s3://<bucket>/<prefix>"), True)
    else:
        print("文件重命名失败")
else:
    print("未找到符合条件的part文件")

额外注意事项

  1. 协议选择:EMR中推荐使用s3a://协议(而非s3://),如果遇到权限或兼容性问题,可将路径改为s3a://<bucket>/<prefix>,同时确保EMR实例角色拥有S3读写权限。
  2. coalesce(1)的风险:生产环境中数据量较大时,coalesce(1)会将所有数据集中到单个Executor,可能引发性能瓶颈或OOM,需根据数据量谨慎使用。
  3. 路径完整性:重命名的目标路径必须是完整的S3路径,否则会默认写入HDFS或本地文件系统,导致文件找不到。

内容的提问来源于stack exchange,提问作者Ronnie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 00:43:24