Spark写入S3报错路径已存在:VW转Spark LDA任务故障排查
解决AWS EMR上Spark任务因输出目录问题失败的排查方案
问题背景
我在AWS EMR上运行Spark任务,需要将S3存储桶中/data/vw目录下的VW LDA格式文件(每行格式为| abc:2 def:1 ghi:3...)转换为Spark LDA库可读取的格式(按冒号后次数重复字符串,如abc abc def ghi ghi ghi),输出至S3的/data/spark目录。
编写的Scala代码
import org.apache.spark.{SparkConf, SparkContext} object Vw2SparkLdaFormatConverter { def repeater(s: String): String = { val ssplit = s.split(':') (ssplit(0) + ' ') * ssplit(1).toInt } def main(args: Array[String]) { val inputPath = args(0) val outputPath = args(1) val conf = new SparkConf().setAppName("FormatConverter") val sc = new SparkContext(conf) val vwdata = sc.textFile(inputPath) val sparkdata = vwdata.map(s => s.trim().split(' ').map(repeater).mkString) val coalescedSparkData = sparkdata.coalesce(100) coalescedSparkData.saveAsTextFile(outputPath) sc.stop() } }
执行命令
spark-submit --class Vw2SparkLdaFormatConverter --deploy-mode cluster --master yarn --conf spark.yarn.submit.waitAppCompletion=true --executor-memory 4g s3a://mybucket/scripts/myscalajar.jar s3a://mybucket/data/vw s3a://mybucket/data/spark
抛出的异常
18/01/20 00:16:28 ERROR ApplicationMaster: User class threw exception: org.apache.hadoop.mapred.FileAlreadyExistsException: Output directory s3a://mybucket/data/spark already exists org.apache.hadoop.mapred.FileAlreadyExistsException: Output directory s3a://mybucket/data/spark already exists at org.apache.hadoop.mapred.FileOutputFormat.checkOutputSpecs(FileOutputFormat.java:131) at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsHadoopDataset$1.apply$mcV$sp(PairRDDFunctions.scala:1119) at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsHadoopDataset$1.apply(PairRDDFunctions.scala:1096) at org.apache.spark.rdd.PairRDDFunctions$$anonfun$saveAsHadoopDataset$1.apply(PairRDDFunctions.scala:1096) at ...
我已尝试更换新的未存在的输出路径(如/data/spark1),但问题仍未解决。作为Scala和Spark新手,我想了解问题所在。
问题排查与解决方案
一、核心问题分析
虽然错误提示是FileAlreadyExistsException,但更换路径后仍报错,说明问题可能不是表面的目录存在,而是以下几个潜在原因:
1. S3路径协议与元数据问题
- Spark在YARN集群模式下,
s3a://、s3://、s3n://协议的元数据可能存在解析差异,EMR默认推荐使用s3://协议,混用协议可能导致路径识别混乱。 - S3的目录删除存在元数据延迟,即使你手动删除了目录,集群可能仍能读取到旧的元数据信息。
2. 代码中的隐性错误
你的repeater函数没有异常处理:如果原始文件中有不符合abc:N格式的行(比如空值、冒号后不是数字),会抛出NumberFormatException,但异常栈可能被外层的目录错误掩盖。
3. 集群权限与任务重试残留
- EMR集群的IAM角色可能没有足够的S3读写权限,导致权限错误被伪装成目录已存在的异常。
- 之前的任务失败后,YARN自动重试可能提前创建了新路径,或者作业提交的缓存导致实际仍写入旧路径。
二、分步解决方案
步骤1:彻底确认目标路径状态
直接登录AWS S3控制台,检查你新指定的路径(比如s3://mybucket/data/spark1)是否真的不存在。如果存在,手动删除整个目录(包括所有子文件和文件夹),等待5分钟后再提交任务,给S3元数据同步留足时间。
步骤2:统一S3协议,修改提交命令
将所有路径替换为EMR默认支持的s3://协议,避免协议解析问题:
spark-submit --class Vw2SparkLdaFormatConverter --deploy-mode cluster --master yarn --conf spark.yarn.submit.waitAppCompletion=true --executor-memory 4g s3://mybucket/scripts/myscalajar.jar s3://mybucket/data/vw s3://mybucket/data/spark_new
步骤3:优化代码,增加容错与日志
给代码添加异常处理和日志输出,方便排查隐性错误:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.log4j.Logger object Vw2SparkLdaFormatConverter { private val logger = Logger.getLogger(getClass.getName) def repeater(s: String): String = { try { val ssplit = s.split(':') if (ssplit.length != 2) { logger.warn(s"无效格式行: $s") "" } else { (ssplit(0) + ' ') * ssplit(1).toInt } } catch { case e: NumberFormatException => logger.warn(s"数字解析失败: $s", e) "" } } def main(args: Array[String]) { if (args.length != 2) { logger.error("请传入输入路径和输出路径两个参数") System.exit(1) } val inputPath = args(0) val outputPath = args(1) val conf = new SparkConf().setAppName("FormatConverter") val sc = new SparkContext(conf) try { val vwdata = sc.textFile(inputPath) val sparkdata = vwdata.map { s => s.trim().split(' ') .filter(_.nonEmpty) // 过滤空字符串 .map(repeater) .mkString.trim() // 去除首尾多余空格 } val coalescedSparkData = sparkdata.coalesce(100) coalescedSparkData.saveAsTextFile(outputPath) } catch { case e: Exception => logger.error("任务执行失败", e) throw e } finally { sc.stop() } } }
步骤4:强制覆盖输出目录
在提交命令中添加强制覆盖输出的配置,绕过目录检查:
spark-submit --class Vw2SparkLdaFormatConverter --deploy-mode cluster --master yarn \ --conf spark.yarn.submit.waitAppCompletion=true \ --conf spark.hadoop.mapreduce.output.fileoutputformat.overwrite=true \ --conf spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=2 \ --executor-memory 4g s3://mybucket/scripts/myscalajar.jar s3://mybucket/data/vw s3://mybucket/data/spark_new
步骤5:检查EMR IAM角色权限
确认EMR集群的服务角色(如EMR_DefaultRole)和实例角色(如EMR_EC2_DefaultRole)拥有以下S3权限:
s3:GetObject(读取输入路径)s3:PutObject、s3:DeleteObject、s3:ListBucket(写入和管理输出路径)
内容的提问来源于stack exchange,提问作者Nik
相关产品推荐
相关产品推荐

