AWS EMR中PySpark脚本作为步骤运行失败求助
排查AWS EMR上PySpark脚本运行失败的问题
嘿,我来帮你一步步揪出问题所在——毕竟本地能正常跑、EMR上翻车,大多是脚本适配或者环境权限的锅,咱们从最容易改的地方开始:
1. 先把脚本里的明显错误修正
你的脚本里有几个不符合EMR Spark运行规范的地方,先改了再说:
- 别直接手动创建
SparkContext!EMR集群会自动帮你初始化好上下文,手动创建大概率会导致冲突。应该先创建SparkSession,再从它里面获取上下文。 sc.read.csv是错的!read方法是属于SparkSession的,不是SparkContext,用spark.read.csv才对。df.select(df)这行完全没意义啊😂,如果要选所有列,直接用df就行;要是选特定列,得写df.select("列名1", "列名2")或者df.select(*df.columns)(选所有列的冗余写法)。- 写入S3的时候,如果目标路径已经存在,Spark默认会报错,记得加个
mode="overwrite"或者mode="append"来覆盖或者追加。
修正后的参考脚本:
from pyspark.sql import SparkSession # 初始化SparkSession,EMR会自动处理集群参数,不用写master spark = SparkSession.builder.appName("S3CSVProcessor").getOrCreate() sc = spark.sparkContext # 读取S3上的CSV文件 df = spark.read.csv("s3://folder1/file.csv", header=True, inferSchema=True) # 如果要选所有列,直接用df就行;选特定列的话改下面这行 dd = df # 示例:dd = df.select("user_id", "order_amount") # 写入S3,加上模式避免路径已存在的报错 write_to = "s3://spark-workflow-test/" dd.write.mode("overwrite").csv(write_to, sep=";", header=True) # 停止会话 spark.stop()
2. 检查EMR的权限配置(大概率是这个问题!)
本地能跑是因为你的本地AWS凭证有S3权限,但EMR集群用的是IAM角色,权限不够的话就会失败:
- 先看EMR的EC2实例角色(默认是
EMR_EC2_DefaultRole),得给它加访问s3://folder1/和s3://spark-workflow-test/的权限:需要包含s3:GetObject(读源文件)、s3:PutObject(写目标文件)、s3:ListBucket(列桶里的内容)这几个权限。 - 还有EMR服务角色(默认
EMR_DefaultRole),确保它有管理集群的完整权限,别被自定义策略限制了。 - 如果用了自定义的EC2实例配置文件,也要确认对应的IAM角色权限到位。
3. EMR环境的其他适配问题
- Spark版本对齐:本地用的Spark版本和EMR上的版本是不是一致?比如本地是Spark 3.3,EMR用的是2.4,有些语法可能不兼容,比如
inferSchema的行为或者写入参数。 - S3路径验证:先在EMR主节点上用
aws s3 ls s3://folder1/file.csv测试下能不能访问到源文件,要是连这个命令都失败,那肯定是路径拼错或者权限问题。 - 网络配置:如果集群在私有子网,得确保有S3访问通道——要么配置了S3的VPC网关端点,要么集群有公网访问权限(如果S3桶不是私有访问的话)。
4. 看日志找具体错误
要是上面的都排查了还不行,就得看EMR的日志了,这是定位问题的关键:
- 登录EMR控制台,找到你的集群,点「步骤」,找到失败的步骤,点「日志链接」,里面会有详细的错误堆栈(比如权限拒绝、文件不存在、内存不足这些都会写得清清楚楚)。
- 也可以直接登录EMR主节点,去
/var/log/spark/目录下找spark-step-*.log文件,里面的信息更全。
5. 其他小概率问题
- 源文件格式问题:要是源CSV有格式错误(比如分隔符不对、特殊字符乱码),
inferSchema=True可能会导致读取失败,你可以先把这个参数去掉,用默认的字符串类型读取,看能不能跑通。 - 内存不够:如果文件很大,EMR实例的内存不足会导致OOM,你可以在添加Spark步骤的时候加配置参数,比如
--conf spark.executor.memory=4g --conf spark.driver.memory=2g来调整内存。
内容的提问来源于stack exchange,提问作者JanBennk
相关产品推荐
相关产品推荐

