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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 04:10:43