【问题标题】:DataFrame orderBy followed by limit in SparkDataFrame orderBy 后跟 Spark 中的限制
【发布时间】:2016-07-28 02:25:01
【问题描述】:

我正在让一个程序生成一个 DataFrame,它将在上面运行类似的东西

    Select Col1, Col2...
    orderBy(ColX) limit(N)

但是,当我最终收集数据时,我发现如果我采用足够大的顶部 N 会导致驱动程序 OOM

另外一个观察是,如果我只做 sort 和 top,这个问题就不会发生。所以这只有在同时存在 sort 和 top 时才会发生。

我想知道为什么会发生这种情况?特别是,在这两种转换组合之下到底发生了什么? spark如何评估排序和限制的查询以及下面对应的执行计划是什么?

还只是好奇 Spark 处理 DataFrame 和 RDD 之间的排序和顶部不同吗?

编辑, 对不起,我不是说收集, 我最初的意思是当我调用任何动作来实现数据时,无论它是否收集(或任何将数据发送回驱动程序的动作)(所以问题绝对不在于输出大小)

【问题讨论】:

  • 我建议您在问题中添加代码。 OOM 错误通常是驱动程序方面的。但是没有真正的代码很难说。

标签: apache-spark


【解决方案1】:

虽然不清楚为什么在这种特殊情况下会失败,但您可能会遇到多个问题:

  • 当您使用limit 时,它只是将所有数据放在一个分区上,无论n 有多大。因此,虽然它没有明确收集它几乎一样糟糕。
  • 除此之外,orderBy 需要使用范围分区进行完全洗牌,当数据分布偏斜时可能会导致不同的问题。
  • 最后当你collect 时,结果可能会大于驱动程序上可用的内存量。

如果你collect 无论如何,这里没有太多可以改进的地方。最终,驱动程序内存将是一个限制因素,但仍有一些可能的改进:

  • 首先不要使用limit
  • collect 替换为toLocalIterator
  • 使用orderBy |> rdd |> zipWithIndex |> filter 或者如果值的确切数量不是硬性要求filter 直接基于近似分布的数据,如Saving a spark dataframe in multiple parts without repartitioning 所示(在 Spark 2.0.0+ 中有方便的approxQuantile 方法)。

【讨论】:

  • 你能写一个如何使用ZipIndex的例子吗?
猜你喜欢
  • 1970-01-01
  • 2020-09-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-07-19
  • 2019-07-04
  • 2020-04-06
  • 2016-03-18
相关资源
最近更新 更多