【问题标题】:Split Spark DataFrame based on condition根据条件拆分 Spark DataFrame
【发布时间】:2017-06-17 00:51:13
【问题描述】:

我需要类似于 randomSplit 函数的东西:

val Array(df1, df2) = myDataFrame.randomSplit(Array(0.6, 0.4))

但是,我需要根据布尔条件拆分 myDataFrame。是否存在类似以下内容?

val Array(df1, df2) = myDataFrame.booleanSplit(col("myColumn") > 100)

我不想做两个单独的 .filter 调用。

【问题讨论】:

    标签: scala apache-spark dataframe apache-spark-sql


    【解决方案1】:

    不幸的是,DataFrame API 没有这样的方法,要按条件拆分,您必须执行两个单独的filter 转换:

    myDataFrame.cache() // recommended to prevent repeating the calculation
    
    val condition = col("myColumn") > 100
    val df1 = myDataFrame.filter(condition)
    val df2 = myDataFrame.filter(not(condition))
    

    【讨论】:

    • 注意:如果此特定示例中的 myColumnNULL,则不会导致正确拆分。您将松开具有 NULL 的列,因为该列不会在 (> 100) 或 (
    【解决方案2】:

    我知道两次缓存和过滤看起来有点难看,但请记住,DataFrame 会被转换为 RDD,它们会被延迟评估,即仅当它们直接或间接用于操作时。

    如果存在问题中建议的方法booleanSplit,则结果将被转换为两个RDD,每个RDD都会被延迟评估。两个 RDD 中的一个将首先评估,另一个将在第二个评估,严格在第一个之后。在评估第一个 RDD 时,第二个 RDD 还没有“出现”(编辑:刚刚注意到 RDD API 有一个类似的问题,answer that gives a similar reasoning

    为了真正获得任何性能优势,第二个 RDD 必须(部分)在第一个 RDD 的迭代期间(或者,实际上,在两者的父 RDD 的迭代期间,这是由第一个 RDD)。 IMO 这与 RDD API 的其余部分的设计不太吻合。不确定性能提升是否可以证明这一点。

    我认为你能做到的最好的方法是避免直接在你的业务代码中编写两个过滤器调用,通过编写一个 implicit class 和一个方法 booleanSplit 作为一个实用方法,它的作用类似于 Tzach Zohar's answer ,也许使用类似于myDataFrame.withColumn("__condition_value", condition).cache() 的东西,因此条件的值不会计算两次。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-14
      • 2021-07-12
      • 1970-01-01
      • 2010-10-31
      • 1970-01-01
      相关资源
      最近更新 更多