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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:52:15