EMR环境下spark.write.parquet写入本地路径及读取方案咨询
EMR上PySpark路径问题解决方案
1. 确保输出写入指定相对路径
Spark写Parquet(HDFS路径)
EMR默认Spark的文件系统是HDFS,配置里的/outputs/features/如果是指向HDFS根目录下的路径,直接使用绝对路径即可:
spark.write.parquet("/outputs/features/")
如果需要基于当前HDFS工作目录的相对路径,先获取当前工作目录再拼接:
from py4j.java_gateway import java_import java_import(spark.sparkContext._jvm, "org.apache.hadoop.fs.FileSystem") fs = spark.sparkContext._jvm.FileSystem.get(spark.sparkContext.hadoopConfiguration) current_hdfs_dir = fs.getWorkingDirectory().toString() output_path = f"{current_hdfs_dir}/outputs/features/" spark.write.parquet(output_path)
Pandas写CSV(本地路径)
df.toPandas().to_csv()是在Driver节点执行的,只会写到Driver节点的本地磁盘。要确保路径存在,先创建目录:
import os local_output_dir = "/outputs/features/" if not os.path.exists(local_output_dir): os.makedirs(local_output_dir) df.toPandas().to_csv(f"{local_output_dir}/result.csv")
注意:EMR本地磁盘是临时存储,集群销毁后数据会丢失,持久化建议优先用HDFS或S3。
2. 从HDFS读取Parquet文件的方法
如果Parquet已经写入HDFS,后续脚本直接用Spark API指定路径即可:
- 读取HDFS绝对路径:
df = spark.read.parquet("/outputs/features/")
- 读取相对于当前用户HDFS home目录的相对路径(比如
/user/hadoop/outputs/features/):
df = spark.read.parquet("outputs/features/")
- 若要先确认文件存在,可在EMR终端执行HDFS命令:
hdfs dfs -ls /outputs/features/
如果之前把Pandas生成的CSV写到了本地,想转到HDFS供后续读取,可执行:
hdfs dfs -put /outputs/features/result.csv /outputs/features/
或者直接用Spark把Pandas DataFrame写入HDFS,避免本地存储:
spark.createDataFrame(df.toPandas()).write.csv("/outputs/features/result.csv")
内容的提问来源于stack exchange,提问作者pnv
相关产品推荐
相关产品推荐

