【问题标题】:Spark job with large text file in gzip format带有 gzip 格式的大文本文件的 Spark 作业
【发布时间】:2016-10-12 03:49:21
【问题描述】:

我正在运行一个 Spark 作业,该作业需要很长时间来处理输入文件。 Gzip 格式的输入文件为 6.8 GB,包含 1.1 亿行文本。我知道它是 Gzip 格式的,所以它是不可拆分的,并且只有一个执行器将用于读取该文件。

作为调试过程的一部分,我决定看看将 gzip 文件转换为 parquet 需要多长时间。我的想法是,一旦我转换为 parquet 文件,然后在该文件上运行我的原始 Spark 作业,在这种情况下,它将使用多个执行程序,并且输入文件将被并行处理。

但即使是很小的工作也比我预期的要花很长时间。这是我的代码:

val input = sqlContext.read.text("input.gz")
input.write.parquet("s3n://temp-output/")

当我在笔记本电脑(16 GB RAM)中提取该文件时,只用了不到 2 分钟。当我在 Spark 集群上运行它时,我的预期是它会花费相同甚至更少的时间,因为我使用的执行器内存是 58 GB。花了大约 20 分钟。

我在这里缺少什么?如果这听起来很业余,我很抱歉,但我在 Spark 中相当新。

在 gzip 文件上运行 Spark 作业的最佳方式是什么?假设我没有选择以其他文件格式(bzip2、snappy、lzo)创建该文件。

【问题讨论】:

  • 您好,您说 parquet 文件处理作业(gzip 到 parquet 需要 20 分钟后)是在驱动程序上执行还是作业提交到集群?您可以通过查看该特定工作的 spark-ui 来检查并判断。如果它在集群上运行,它将显示集群上的多个节点。
  • 不是,是在集群上提交的。

标签: hadoop apache-spark amazon-s3 spark-dataframe parquet


【解决方案1】:

在进行输入-处理-输出类型的 Spark 作业时,需要考虑三个不同的问题:

  1. 输入并行度
  2. 处理并行度
  3. 输出并行度

在您的情况下,输入并行度为 1,因为在您的问题中您声称无法更改输入格式或粒度。

您也基本上没有进行任何处理,因此您无法在那里获得任何收益。

但是,您可以控制输出并行度,这将为您带来两个好处:

  • 多个 CPU 会写入,从而减少写入操作的总时间。

  • 您的输出将被拆分为多个文件,以便您在以后的处理中利用输入并行性。

为了增加并行度,你必须增加分区的数量,你可以用repartition()来做,例如,

val numPartitions = ...
input.repartition(numPartitions).write.parquet("s3n://temp-output/")

在选择最佳分区数时,需要考虑许多不同的因素。

  • 数据大小
  • 分区倾斜
  • 集群 RAM 大小
  • 集群中的核心数
  • 您将执行的后续处理类型
  • 您将用于后续处理的集群大小(RAM 和内核)
  • 您正在写入的系统

在不了解您的目标和限制的情况下,很难做出可靠的建议,但这里有一些通用的指导原则:

  • 由于您的分区不会倾斜(上面对repartition 的使用将使用一个哈希分区器来纠正倾斜),如果您将分区数设置为等于数字,您将获得最快的吞吐量执行器核心,假设您使用的节点具有足够的 I/O。

  • 当您处理数据时,您确实希望整个分区能够“适应”分配给单个执行程序核心的 RAM。 “适合”在这里的含义取决于您的处理。如果您正在执行简单的map 转换,则数据可能会被流式传输。如果您正在做一些涉及订购的事情,那么 RAM 需求就会大幅增长。如果您使用的是 Spark 1.6+,您将受益于更灵活的内存管理。如果您使用的是早期版本,则必须更加小心。当 Spark 必须开始“缓冲”到磁盘时,作业执行会停止。磁盘上的大小和内存中的大小可能非常非常不同。后者取决于您处理数据的方式以及 Spark 可以从谓词下推中获得多少好处(Parquet 支持这一点)。使用 Spark UI 查看各个作业阶段占用多少 RAM。

顺便说一句,除非您的数据具有非常特定的结构,否则不要对分区号进行硬编码,因为这样您的代码将在不同大小的集群上以次优方式运行。相反,请使用以下技巧来确定集群中的执行程序数量。然后,您可以根据您使用的机器乘以每个执行程序的核心数。

// -1 is for the driver node
val numExecutors = sparkContext.getExecutorStorageStatus.length - 1

作为参考,在我们的团队中,我们使用相当复杂的数据结构,这意味着 RAM 大小 >> 磁盘大小,我们的目标是将 S3 对象保持在 50-250Mb 范围内,以便在每个节点上进行处理执行器核心有 10-20Gb RAM。

希望这会有所帮助。

【讨论】:

  • 感谢@Sim 提供的详细信息,这真的很有帮助。我现在正在根据您的评论尝试不同的设置。仅供参考:我使用 3 台 r3.4xlarge 机器作为 EMR 集群的核心。所以每个节点有 16 个 vCPU,122 GB 内存。 Spark 版本是 1.6.1。我在关注这个:blog.cloudera.com/blog/2015/03/… 调整# of executors/cores/memory 但没有得到预期的结果。当我尝试不同的设置时,您会根据我刚刚分享的信息推荐任何其他设置吗?再次感谢。
  • 你将如何处理持久化的数据?
  • 我需要读取持久化的数据(有 110M 行的字符串)并且需要执行两个 Flapmap 来创建一些值对,稍后我需要执行 reduceByKey 以通过键对聚合它们。最后,我需要将数据映射到其他格式并计算一些其他值(联合、相交),最后将它们保存到 Redshift 中。它正在工作,但正在尝试优化它。
  • 重新分区有助于将时间从 23 分钟缩短到 15/12 分钟,其中读取耗时 7.9 分钟,写入耗时 3.5 分钟(60 个分区)或 2 分钟(120 个分区)。
  • @dreamer 我很高兴这有帮助!根据您描述的处理需求,您在分区方面有很大的灵活性。只需保持分区数 >= 内核数即可。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2015-03-02
  • 2017-03-22
  • 2011-01-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-05-19
相关资源
最近更新 更多