【问题标题】:Repartition with Apache Spark使用 Apache Spark 重新分区
【发布时间】:2016-11-14 10:32:54
【问题描述】:

问题:我正在尝试对数据集进行重新分区,以便指定整数列中具有相同编号的所有行都位于同一分区中。

工作原理:当我将 1.6 API(Java 中)与 RDD 一起使用时,我使用了一个哈希分区器,它按预期工作。例如,如果我为每一行打印此列的每个值的模数,我在给定分区中得到相同的模数(我通过手动读取使用 saveAsHadoopFile 保存的内容来读取分区)。

使用最新的 API 无法正常工作

但现在我正在尝试使用 2.0.1 API(在 Scala 中)和具有重新分区方法的数据集,该方法采用多个分区和一列并将此数据集保存为镶木地板文件。如果我在给定此列的情况下查看未分区行的分区,结果会有所不同。

【问题讨论】:

    标签: java scala hadoop apache-spark


    【解决方案1】:

    要保存已分区的Dataset,您可以使用:

    • DataFrameWriter.partitionBy - 从 Spark 1.6 开始可用

      df.write.partitionBy("someColumn").format(...).save()
      
    • DataFrameWriter.bucketBy - 从 Spark 2.0 开始可用

      df.write.bucketBy("someColumn").format(...).save()
      

    使用df.partitionBy("someColumn").write.format(...).save 应该也可以,但Dataset API 不使用哈希码。它使用MurmurHash,因此结果将与RDD API 中HashParitioner 的结果不同,并且琐碎的检查(如您描述的那样)将不起作用。

    val oldHashCode = udf((x: Long) => x.hashCode)
    
    // https://github.com/apache/spark/blob/v2.0.1/core/src/main/scala/org/apache/spark/util/Utils.scala#L1596-L1599
    val nonNegativeMode = udf((x: Int, mod: Int) => {
      val rawMod = x % mod
      rawMod + (if (rawMod < 0) mod else 0)
    })
    
    val df = spark.range(0, 10)
    
    val oldPart = nonNegativeMode(oldHashCode($"id"), lit(3))
    val newPart = nonNegativeMode(hash($"id"), lit(3))
    
    df.select($"*", oldPart, newPart).show
    
    +---+---------------+--------------------+
    | id|UDF(UDF(id), 3)|UDF(hash(id, 42), 3)|
    +---+---------------+--------------------+
    |  0|              0|                   1|
    |  1|              1|                   2|
    |  2|              2|                   2|
    |  3|              0|                   0|
    |  4|              1|                   2|
    |  5|              2|                   2|
    |  6|              0|                   0|
    |  7|              1|                   0|
    |  8|              2|                   2|
    |  9|              0|                   2|
    +---+---------------+--------------------+
    

    一个可能的问题是DataFrame writer 可以合并多个小文件以降低成本,因此可以将来自不同分区的数据放在一个文件中。

    【讨论】:

    • 谢谢!我刚刚用 bucketBy 测试了你的示例,因为它似乎完全符合我的要求(同一分区中给定列中具有相同数字的行)但在执行的某个时刻我得到Exception in thread "main" org.apache.spark.sql.AnalysisException: 'save' does not support bucketing right now;
    • 我有类似的问题 (stackoverflow.com/questions/49434262/…) 但我想用一些自定义逻辑编写数据。因此,不能直接使用 write 方法。你能帮忙吗?
    • 我遇到了类似的问题,我正在加载一个数据框,对其进行分区并按列排序,然后使用 save() 写入它。我得到同样的错误“org.apache.spark.sql.AnalysisException:'save'现在不支持分桶”。你能解释一下你是如何解决这个问题的吗?
    猜你喜欢
    • 2017-04-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-02
    • 2022-08-03
    • 2019-08-26
    • 2015-10-15
    • 1970-01-01
    相关资源
    最近更新 更多