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文件")
额外注意事项
- 协议选择:EMR中推荐使用
s3a://协议(而非s3://),如果遇到权限或兼容性问题,可将路径改为s3a://<bucket>/<prefix>,同时确保EMR实例角色拥有S3读写权限。 - coalesce(1)的风险:生产环境中数据量较大时,
coalesce(1)会将所有数据集中到单个Executor,可能引发性能瓶颈或OOM,需根据数据量谨慎使用。 - 路径完整性:重命名的目标路径必须是完整的S3路径,否则会默认写入HDFS或本地文件系统,导致文件找不到。
内容的提问来源于stack exchange,提问作者Ronnie
相关产品推荐
相关产品推荐

