【问题标题】:Partition data for efficient joining for Spark dataframe/dataset对 Spark 数据帧/数据集进行有效连接的分区数据
【发布时间】:2018-06-18 01:00:06
【问题描述】:

我需要join 基于一些共享键列的多个 DataFrame。对于键值 RDD,可以指定一个分区器,以便将具有相同键的数据点洗牌到同一个执行器,这样加入效率更高(如果在 join 之前有洗牌相关操作)。可以在 Spark DataFrames 或 DataSets 上做同样的事情吗?

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-dataframe partitioning apache-spark-dataset


    【解决方案1】:

    如果你知道你会多次加入,你可以在加载后repartition一个DataFrame

    val users = spark.read.load("/path/to/users").repartition('userId)
    
    val joined1 = users.join(addresses, "userId")
    joined1.show() // <-- 1st shuffle for repartition
    
    val joined2 = users.join(salary, "userId")
    joined2.show() // <-- skips shuffle for users since it's already been repartitioned
    

    所以它会洗牌一次,然后在以后加入时重复使用洗牌文件。

    但是,如果您知道自己会在某些键上反复打乱数据,那么最好的办法是将数据保存为分桶表。这会将数据写出已经预先散列分区的数据,因此当您读取表并加入它们时,您可以避免洗牌。你可以这样做:

    // you need to pick a number of buckets that makes sense for your data
    users.bucketBy(50, "userId").saveAsTable("users")
    addresses.bucketBy(50, "userId").saveAsTable("addresses")
    
    val users = spark.read.table("users")
    val addresses = spark.read.table("addresses")
    
    val joined = users.join(addresses, "userId")
    joined.show() // <-- no shuffle since tables are co-partitioned
    

    为了避免洗牌,表必须使用相同的分桶(例如,相同数量的桶和桶列上的连接)。

    【讨论】:

    • 单连接使用repartition有好处吗?
    • 它不会伤害但不会提供任何好处。 Spark 无论如何都需要洗牌才能加入,除非它是广播。
    • 我试图理解为什么 Spark 要优化第二个作业,您需要通过 userId 显式重新分区。在第一个需要随机播放的作业之后,Spark 不会知道数据现在按用户 ID 分区吗?
    【解决方案2】:

    可以通过 repartition 方法使用 DataFrame/DataSet API。使用此方法,您可以指定一个或多个列用于数据分区,例如

    val df2 = df.repartition($"colA", $"colB")
    

    也可以在同一命令中同时指定想要的分区数,

    val df2 = df.repartition(10, $"colA", $"colB")
    

    注意:这并不能保证数据帧的分区将位于同一个节点上,只保证分区以相同的方式完成。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-01-23
      • 2017-08-16
      • 2017-09-16
      相关资源
      最近更新 更多