【问题标题】:Multipart upload error to S3 from Spark从 Spark 到 S3 的分段上传错误
【发布时间】:2017-02-05 04:00:07
【问题描述】:

我在尝试关闭序列文件编写器时收到错误消息“第 2 部分的上传尝试已达到最大限制:5,将引发异常并失败”。完整的异常日志如下:

16/12/30 19:47:01 INFO s3n.MultipartUploadOutputStream: uploadPart /mnt/s3/57b63810-c20a-438c-a73f-48d50e0be7d2-0001 94317523 bytes md5: 05ww/fe3pNni9Zvfm+l4Gg== md5hex: d39c30fdf7b7a4d9e2f59bdf9be9781a 16/12/30 19:47:12 INFO s3n.MultipartUploadOutputStream: uploadPart /mnt1/s3/57b63810-c20a-438c-a73f-48d50e0be7d2-0002 94317523 bytes md5: 05ww/fe3pNni9Zvfm+l4Gg== md5hex: d39c30fdf7b7a4d9e2f59bdf9be9781a 16/12/30 19:47:23 INFO s3n.MultipartUploadOutputStream: uploadPart /mnt/s3/57b63810-c20a-438c-a73f-48d50e0be7d2-0003 94317523 bytes md5: 05ww/fe3pNni9Zvfm+l4Gg== md5hex: d39c30fdf7b7a4d9e2f59bdf9be9781a 16/12/30 19:47:35 INFO s3n.MultipartUploadOutputStream: uploadPart /mnt1/s3/57b63810-c20a-438c-a73f-48d50e0be7d2-0004 94317523 bytes md5: 05ww/fe3pNni9Zvfm+l4Gg== md5hex: d39c30fdf7b7a4d9e2f59bdf9be9781a 16/12/30 19:47:46 INFO s3n.MultipartUploadOutputStream: uploadPart /mnt/s3/57b63810-c20a-438c-a73f-48d50e0be7d2-0005 94317523 bytes md5: 05ww/fe3pNni9Zvfm+l4Gg== md5hex: d39c30fdf7b7a4d9e2f59bdf9be9781a 30 年 16 月 12 日 19:47:57 错误 s3n.MultipartUploadOutputStream:第 2 部分的上传尝试已达到最大限制:5,将引发异常并失败 16/12/30 19:47:57 信息 s3n.MultipartUploadOutputStream:key 的 completeMultipartUpload 错误:输出/part-20176 java.lang.IllegalStateException:达到部分上传尝试的最大限制 在 com.amazon.ws.emr.hadoop.fs.s3n.MultipartUploadOutputStream.spawnNewFutureIfNeeded(MultipartUploadOutputStream.java:362) 在 com.amazon.ws.emr.hadoop.fs.s3n.MultipartUploadOutputStream.uploadMultiParts(MultipartUploadOutputStream.java:422) 在 com.amazon.ws.emr.hadoop.fs.s3n.MultipartUploadOutputStream.close(MultipartUploadOutputStream.java:471) 在 org.apache.hadoop.fs.FSDataOutputStream$PositionCache.close(FSDataOutputStream.java:74) 在 org.apache.hadoop.fs.FSDataOutputStream.close(FSDataOutputStream.java:108) 在 org.apache.hadoop.io.SequenceFile$Writer.close(SequenceFile.java:1290) ... 在 org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$18.apply(RDD.scala:727) 在 org.apache.spark.rdd.RDD$$anonfun$mapPartitionsWithIndex$1$$anonfun$apply$18.apply(RDD.scala:727) 在 org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38) 在 org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:300) 在 org.apache.spark.rdd.RDD.iterator(RDD.scala:264) 在 org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:66) 在 org.apache.spark.scheduler.Task.run(Task.scala:88) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:214) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) 在 java.lang.Thread.run(Thread.java:745) 16/12/30 19:47:59 信息 s3n.MultipartUploadOutputStream:uploadPart 错误 com.amazonaws.AbortedException: 16/12/30 19:48:18 信息 s3n.MultipartUploadOutputStream:uploadPart 错误 com.amazonaws.AbortedException:

我刚刚收到 5 次重试失败的错误。我不明白为什么。有没有人见过这个错误?这可能是什么原因?

我正在使用我自己的多输出格式实现来编写序列文件:

class MultiOutputSequenceFileWriter(prefix: String, suffix: String) extends Serializable {
   private val writers = collection.mutable.Map[String, SequenceFile.Writer]()

   /**
     * @param pathKey    folder within prefix where the content will be written
     * @param valueKey   key of the data to be written
     * @param valueValue value of the data to be written
     */
   def write(pathKey: String, valueKey: Any, valueValue: Any) = {
     if (!writers.contains(pathKey)) {
       val path = new Path(prefix + "/" + pathKey + "/" + "part-" + suffix)
       val hadoopConf = new conf.Configuration()
       hadoopConf.setEnum("io.seqfile.compression.type", SequenceFile.CompressionType.NONE)
       val fs = FileSystem.get(hadoopConf)
       writers(pathKey) = SequenceFile.createWriter(hadoopConf, Writer.file(path),
         Writer.keyClass(valueKey.getClass()),
         Writer.valueClass(valueValue.getClass()),
         Writer.bufferSize(fs.getConf().getInt("io.file.buffer.size", 4096)), //4KB
         Writer.replication(fs.getDefaultReplication()),
         Writer.blockSize(1073741824), // 1GB
         Writer.progressable(null),
         Writer.metadata(new Metadata()))
     }
     writers(pathKey).append(valueKey, valueValue)
   }
   def close = writers.values.foreach(_.close())
}

我正在尝试将序列文件编写如下:

...
rdd.mapPartitionsWithIndex { (p, it) => {
  val writer = new MultiOutputSequenceFileWriter("s3://bucket/output/", p.toString)
  for ( (key1, key2, data) <- it) {
    ...
    writer.write(key1, key2, data)
    ...
  }
  writer.close
  Nil.iterator
}.foreach( (x:Nothing) => ()) // To trigger iterator
}
...

注意:

  • 当我尝试关闭编写器时遇到异常(我认为编写器尝试在关闭之前编写内容,我认为异常是由于这个原因出现的)。
  • 我用相同的输入再试了两次相同的作业。我在第一次重新运行时没有出错,但在第二次中出现了三个错误。这可能只是 S3 中的暂时性问题吗?
  • 失败的零件文件在 S3 中不存在。

【问题讨论】:

    标签: apache-spark amazon-s3 emr


    【解决方案1】:

    AWS 支持工程师提到,在发生错误时,存储桶上有很多命中。该作业正在重试默认次数 (5),并且很可能所有重试都被限制了。现在,我在提交作业时添加了以下配置参数来增加重试次数。

    spark.hadoop.fs.s3.maxRetries=20

    此外,我还压缩了输出,以便减少对 S3 的请求数量。在这些更改之后,我没有看到多次运行失败。

    【讨论】:

      【解决方案2】:

      编写器(亚马逊代码顺便说一句,spark 或 hadoop 团队不会处理任何事情)在数据生成时(在后台线程中)将数据写入块中,其余数据和多部分上传在 close() 中提交——这也是代码将阻止等待所有挂起的上传完成的地方。

      听起来有些 PUT 请求已经失败,并且在 close() 调用中会发现并报告此失败。我不知道 EMR s3:// 客户端是否使用该块大小作为其分区的大小标记;我个人会推荐一个较小的尺寸,比如 128MB。

      无论如何:假设暂时的网络问题,或者您分配的 EC2 虚拟机的网络连接不良。要求一个新的虚拟机。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2013-03-29
        • 2016-03-23
        • 2015-12-28
        • 1970-01-01
        • 2017-08-23
        • 2014-10-14
        • 2014-04-30
        相关资源
        最近更新 更多