【问题标题】:PySpark Yarn Application fails on groupByPySpark 纱线应用程序在 groupBy 上失败
【发布时间】:2016-02-05 06:21:00
【问题描述】:

我正在尝试在 Yarn 模式下运行一项作业,该作业处理从谷歌云存储读取的大量数据 (2TB)。

管道可以这样总结:

sc.textFile("gs://path/*.json")\
.map(lambda row: json.loads(row))\
.map(toKvPair)\
.groupByKey().take(10)

 [...] later processing on collections and output to GCS.
  This computation over the elements of collections is not associative,
  each element is sorted in it's keyspace.

10GB 数据上运行时,它已完成,没有任何问题。 但是,当我在完整数据集上运行它时,它总是在容器中出现此日志时失败:

15/11/04 16:08:07 WARN org.apache.spark.scheduler.cluster.YarnSchedulerBackend$YarnSchedulerEndpoint: ApplicationMaster has disassociated: xxxxxxxxxxx
15/11/04 16:08:07 ERROR org.apache.spark.scheduler.cluster.YarnClientSchedulerBackend: Yarn application has already exited with state FINISHED!
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/usr/lib/spark/python/pyspark/rdd.py", line 1299, in take
    res = self.context.runJob(self, takeUpToNumLeft, p)
  File "/usr/lib/spark/python/pyspark/context.py", line 916, in runJob
15/11/04 16:08:07 WARN org.apache.spark.ExecutorAllocationManager: No stages are running, but numRunningTasks != 0
    port = self._jvm.PythonRDD.runJob(self._jsc.sc(), mappedRDD._jrdd, partitions)
  File "/usr/lib/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/java_gateway.py", line 538, in __call__
  File "/usr/lib/spark/python/pyspark/sql/utils.py", line 36, in deco
    return f(*a, **kw)
  File "/usr/lib/spark/python/lib/py4j-0.8.2.1-src.zip/py4j/protocol.py", line 300, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.runJob.
: org.apache.spark.SparkException: Job cancelled because SparkContext was shut down

我尝试通过连接到主服务器来逐个启动每个操作进行调查,但它似乎在 groupBy 上失败了。我还尝试通过添加节点并升级它们的内存和 CPU 数量来重新扩展集群,但我仍然遇到同样的问题。

120 个节点 + 1 个具有相同规格的主节点: 8 个 vCPU - 52GB 内存

我试图找到有类似问题的线程但没有成功,所以我真的不知道我应该提供什么信息,因为日志不是很清楚,所以请随时询问更多信息。

主键是每条记录的必填值,我们需要所有没有过滤器的键,大约代表 600k 个键。 真的可以在不将集群扩展到大规模的情况下执行此操作吗?我刚刚读到 databricks 对 100TB 的数据 (https://databricks.com/blog/2014/10/10/spark-petabyte-sort.html) 进行了排序,这也涉及到大规模的洗牌。他们成功地将多个内存缓冲区替换为单个缓冲区,从而导致大量磁盘 IO ?我的集群规模是否可以执行此类操作?

【问题讨论】:

  • 你能分享更多关于你的 toKvPair 功能的细节吗?分组时,每个键大约需要多少个值?它是否会为某些可能没有预期键的记录生成“空”键?
  • 是的,当然,我只是编辑以添加更多信息。
  • 因此,问题的根源可能是 groupByKey 引入了不可扩展的每个键瓶颈,而不是遇到“总数据集大小”约束。在 petasort 基准测试中,通常键本身仍然是唯一的或仅限于很少的重复项,而不管整体数据集大小如何,这使得可以使用更大的集群进行扩展。请参阅 Databricks 的这篇 avoid groupByKey 帖子,了解“热键”如何可能是一个陷阱。
  • 现在,如果 2TB 完全均匀地分布在 600k 个密钥中,这应该不是问题,并且应该可以找到您遇到的其他瓶颈,因为这分为每组只有 3MB。问题在于,如果您有任何倾斜的高度流行的密钥,您需要将 4GB 的数据全部放入一个密钥中。然后 groupByKey 会遇到问题。
  • 您可以尝试使用countByKey 而不是 groupByKey,然后对它们进行排序并获得最大计数吗?即使您有倾斜的键,计数也应该是可行的,因为它应该能够增量地合并求和,而无需将单个键的所有值加载到内存中。

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


【解决方案1】:

总结一下我们在原始问题上通过 cmets 学到的知识,如果一个小数据集有效(尤其是一个可能适合单台机器总内存的数据集),然后一个大数据集失败,尽管向集群添加了更多的节点,结合对于groupByKey 的任何用法,最常见的问题是您的数据是否存在每个键的记录数显着不平衡。

特别是,groupByKey 直到今天仍然有一个限制,即不仅单个键的所有值都必须洗牌到同一台机器上,它们还必须能够适应内存:

https://github.com/apache/spark/blob/master/core/src/main/scala/org/apache/spark/rdd/PairRDDFunctions.scala#L503

/**
 * Group the values for each key in the RDD into a single sequence. Allows controlling the
 * partitioning of the resulting key-value pair RDD by passing a Partitioner.
 * The ordering of elements within each group is not guaranteed, and may even differ
 * each time the resulting RDD is evaluated.
 *
 * Note: This operation may be very expensive. If you are grouping in order to perform an
 * aggregation (such as a sum or average) over each key, using [[PairRDDFunctions.aggregateByKey]]
 * or [[PairRDDFunctions.reduceByKey]] will provide much better performance.
 *
 * Note: As currently implemented, groupByKey must be able to hold all the key-value pairs for any
 * key in memory. If a key has too many values, it can result in an [[OutOfMemoryError]].
 */

有一些further discussion of this problem 指向mailing list discussion,其中包括一些解决方法的讨论;也就是说,您可以将值/记录的哈希字符串显式附加到键中,散列到一些小的存储桶集中,以便您手动分片您的大组。

在您的情况下,您甚至可以最初进行.map 转换,它只会有条件地调整已知热键的键以将其划分为子组,同时保持非热键不变。

一般来说,“内存中”约束意味着您无法通过添加更多节点来真正解决显着倾斜的键,因为它需要在热节点上“就地”缩放。对于特定情况,您可以将spark.executor.memory 设置为--conf 或在dataproc 中gcloud beta dataproc jobs submit spark [other flags] --properties spark.executor.memory=30g,只要最大键的值都可以适合30g(还有一些净空/开销)。但这将在任何可用的最大机器上达到顶峰,因此,如果在整个数据集增长时最大密钥的大小可能会增长,最好更改密钥分布本身而不是尝试增加单执行器内存.

【讨论】:

  • 非常感谢您的详细回答!我将键拆分为碎片,将记录时间戳的日期附加到键上,并且它起作用了,这也使数据排序更快,更容易减少,以获得每个键的完整信息。但是,我的管道中仍然会出现“内存不足错误”(在减少键的分片数据之前),但我会在另一篇文章中询问。无论如何,感谢您提供的这些重要提示。
  • 这是问题:stackoverflow.com/questions/33547649/…,我不确定使用一个执行器是否与 groupBy 相关联,但也许您可以提供帮助。
猜你喜欢
  • 2021-11-20
  • 2021-05-27
  • 1970-01-01
  • 1970-01-01
  • 2023-04-02
  • 1970-01-01
  • 1970-01-01
  • 2022-06-14
  • 2021-09-01
相关资源
最近更新 更多