【问题标题】:How to (equally) partition array-data in spark dataframe如何(平等地)在火花数据框中划分数组数据
【发布时间】:2017-09-15 13:24:30
【问题描述】:

我有一个如下形式的数据框:

import scala.util.Random
val localData = (1 to 100).map(i => (i,Seq.fill(Math.abs(Random.nextGaussian()*100).toInt)(Random.nextDouble)))
val df = sc.parallelize(localData).toDF("id","data")

|-- id: integer (nullable = false)
|-- data: array (nullable = true)
|    |-- element: double (containsNull = false)


df.withColumn("data_size",size($"data")).show

+---+--------------------+---------+
| id|                data|data_size|
+---+--------------------+---------+
|  1|[0.77845301260182...|      217|
|  2|[0.28806915178410...|      202|
|  3|[0.76304121847720...|      165|
|  4|[0.57955190088558...|        9|
|  5|[0.82134215959459...|       11|
|  6|[0.42193739241567...|       57|
|  7|[0.76381645621403...|        4|
|  8|[0.56507523859466...|       93|
|  9|[0.83541853717244...|      107|
| 10|[0.77955626749231...|      111|
| 11|[0.83721643562080...|      223|
| 12|[0.30546029947285...|      116|
| 13|[0.02705462199952...|       46|
| 14|[0.46646815407673...|       41|
| 15|[0.66312488908446...|       16|
| 16|[0.72644646115640...|      166|
| 17|[0.32210572380128...|      197|
| 18|[0.66680355567329...|       61|
| 19|[0.87055594653295...|       55|
| 20|[0.96600507545438...|       89|
+---+--------------------+---------+

现在我想应用一个昂贵的 UDF,计算时间与数据数组的大小成正比。我想知道如何重新分区我的数据,以便每个分区具有大致相同数量的“records*data_size”(即,数据点不仅仅是记录)。

如果只是做df.repartition(100),我可能会得到一些包含一些非常大的数组的分区,这些数组是整个 spark 阶段的瓶颈(所有其他任务都已经完成)。如果我当然可以选择数量惊人的分区,这将(几乎)确保每条记录都在一个单独的分区中。但是还有其他方法吗?

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    如您所说,您可以增加分区数量。我通常使用核心数的倍数:spark context default parallelism * 2-3..
    在您的情况下,您可以使用更大的乘数。

    另一种解决方案是以这种方式过滤拆分您的 df:

    • df 只有更大的数组
    • df 与其余部分

    然后您可以对它们中的每一个进行重新分区,执行计算并将它们合并回来。

    请注意,重新分区可能会很昂贵,因为您有大行要随机播放。

    您可以查看这些幻灯片 (27+):https://www.slideshare.net/SparkSummit/custom-applications-with-sparks-rdd-spark-summit-east-talk-by-tejas-patil

    他们遇到了非常严重的数据偏差,不得不以一种有趣的方式处理它。

    【讨论】:

    • 在我的情况下,将数据帧拆分为小/大记录的想法可能就足够了。
    猜你喜欢
    • 2019-04-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-23
    • 1970-01-01
    • 2020-04-21
    相关资源
    最近更新 更多