【问题标题】:Why does Spark running in Google Dataproc store temporary files on external storage (GCS) instead of local disk or HDFS while using saveAsTextFile?为什么在使用 saveAsTextFile 时,在 Google Dataproc 中运行的 Spark 将临时文件存储在外部存储 (GCS) 而不是本地磁盘或 HDFS 上?
【发布时间】:2016-12-15 18:25:33
【问题描述】:

我已经运行了以下 PySpark 代码:

from pyspark import SparkContext

sc = SparkContext()

data = sc.textFile('gs://bucket-name/input_blob_path')
sorted_data = data.sortBy(lambda x: sort_criteria(x))
sorted_data.saveAsTextFile(
    'gs://bucket-name/output_blob_path',
    compressionCodecClass="org.apache.hadoop.io.compress.GzipCodec"
)

作业成功完成。但是,在作业执行期间,Spark 在以下路径 gs://bucket-name/output_blob_path/_temporary/0/ 中创建了许多临时 blob。我意识到最后删除所有这些临时 blob 占用了一半的作业执行时间,而这段时间内 CPU 利用率为 1%(极大浪费资源)。

有没有办法将临时文件存储在本地驱动器(或 HDFS)而不是 GCP 上?我仍然希望将最终结果(排序数据集)保存到 GCP。

我们使用具有 10 个工作节点的 Dataproc Spark 集群(VM 类型 16 核,60GM)。输入数据量为10TB。

【问题讨论】:

    标签: apache-spark pyspark google-cloud-dataproc


    【解决方案1】:

    您看到的_temporary 文件很可能是在后台使用的FileOutputCommitter 的产物。重要的是,这些临时 blob 并不是严格意义上的“临时”数据,而是实际上完成的输出数据,只有在作业完成时才会“重命名”到最终目的地。通过重命名这些文件的“提交”实际上很快,因为源和目标都在 GCS 上;出于这个原因,无法通过将临时文件放在 HDFS 上然后“提交”到 GCS 来替换工作流的那一部分,因为这样提交将需要将整个输出数据集从 HDFS 重新连接到 GCS。具体来说,底层的 Hadoop FileOutputFormat 类不支持这样的习惯用法。

    GCS 本身并不是一个真正的文件系统,而是一个“对象存储”,Dataproc 内部的 GCS 连接器只是尽可能地模仿 HDFS。一个后果是删除目录填充文件实际上需要 GCS 在后台删除单个对象,而不是真正的文件系统只是取消链接 inode。

    在实践中,如果您遇到此问题,则可能意味着您的输出无论如何都被拆分为太多文件,因为清理确实一次以大约 1000 个文件的批次进行。因此,多达数万个输出文件通常不会很慢。拥有太多文件也会使将来处理这些文件的速度变慢。最简单的解决方法通常是尽可能减少输出文件的数量,例如使用repartition()

    from pyspark import SparkContext
    
    sc = SparkContext()
    
    data = sc.textFile('gs://bucket-name/input_blob_path')
    sorted_data = data.sortBy(lambda x: sort_criteria(x))
    sorted_data.repartition(1000).saveAsTextFile(
        'gs://bucket-name/output_blob_path',
        compressionCodecClass="org.apache.hadoop.io.compress.GzipCodec"
    )
    

    【讨论】:

    • 感谢您的解释。我有点惊讶,因为我们对从 BigQuery 导出到 GCS 的数据进行了排序,所以文件太多了。我的假设是 BiqQuery 导出功能已经优化了分区数量(在 GCS 上存储数据集的最佳文件数量)。
    • 根据所应用的 RDD 操作的种类,转换后的分区数量可能与输入分区的数量不同,并且在这种情况下 FileInputFormat 将默认切分输入无论如何,将文件放入较小的分区中,而与输入文件的数量无关。例如,您可以使用 --properties spark.hadoop.fs.gs.block.size=536870912 将其调整为 512MB,而不是默认的 64MB。
    • 您可能还想在集群部署时默认调整它。如果您的工作通常在 10TB 范围内,gcloud dataproc clusters create my-cluster --properties core:fs.gs.block.size=536870912 将是合理的。但是,如果您的工作只有 10GB,那就太高了。在大多数情况下,目标是超过 1000 个或少于 50000 个分区是很好的,但即使是小型作业,通常也不希望使用小于 64MB 的块大小。
    • 非常有用的反馈!谢谢。
    • 我已经回答了 AWS stackoverflow.com/a/54350777/1931239 的类似场景。
    【解决方案2】:

    我和你之前有同样的问题。 My blog: spark speedup write file to cloud storage。然后我找到这篇文章Spark 2.0.0 Cluster Takes a Longer Time to Append Data

    如果您发现使用 Spark 2.0.0 版本的集群需要较长时间将数据追加到现有数据集,特别是所有 Spark 作业都已完成,但您的命令尚未完成,这是因为驱动节点是将任务的输出文件从作业临时目录一个接一个地移动到最终目的地,这对于云存储来说很慢。要解决此问题,请将 mapreduce.fileoutputcommitter.algorithm.version 设置为 2。请注意,此问题不会影响覆盖数据集或将数据写入新位置。

    当我们使用 GCS 作为 tmp 存储时,这个问题在云环境中会被放大。

    如何解决?

    您可以简单地添加此参数来解决此问题,这意味着您在将文件保存到 GCS 时不会创建 tmp 文件。

    write.option("mapreduce.fileoutputcommitter.algorithm.version", "2")
    

    警告!

    由于数据丢失的可能性,DirectParquetOutputCommitter 已从 Spark 2.0 中删除。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2012-01-20
      • 1970-01-01
      • 1970-01-01
      • 2016-09-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-02-05
      相关资源
      最近更新 更多