【问题标题】:How to efficiently identify records that are different for a specific column如何有效地识别特定列的不同记录
【发布时间】:2020-11-15 11:23:11
【问题描述】:

我有两个数据集df1df2,我需要在其中检测df2df1 中不同的任何记录,并创建一个结果数据集,其中包含一个标记不同记录的附加列。这是一个例子。

package playground

import org.apache.log4j.{Level, Logger}
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, sum}

object sample4 {

  val spark = SparkSession
    .builder()
    .appName("Sample app")
    .master("local")
    .getOrCreate()

  val sc = spark.sparkContext

  final case class Owner(a: Long,
                         b: String,
                         c: Long,
                         d: Short,
                         e: String,
                         f: String,
                         o_qtty: Double)

  final case class Result(a: Long,
                          b: String,
                          c: Long,
                          d: Short,
                          e: String,
                          f: String,
                          o_qtty: Double,
                          isDiff: Boolean)

  def main(args: Array[String]): Unit = {
    Logger.getLogger("org").setLevel(Level.OFF)

    import spark.implicits._

    val data1 = Seq(
      Owner(11, "A", 666, 2017, "x", "y", 50),
      Owner(11, "A", 222, 2018, "x", "y", 20),
      Owner(33, "C", 444, 2018, "x", "y", 20),
      Owner(33, "C", 555, 2018, "x", "y", 120),
      Owner(22, "B", 555, 2018, "x", "y", 20),
      Owner(99, "D", 888, 2018, "x", "y", 100),
      Owner(11, "A", 888, 2018, "x", "y", 100),
      Owner(11, "A", 666, 2018, "x", "y", 80),
      Owner(33, "C", 666, 2018, "x", "y", 80),
      Owner(11, "A", 444, 2018, "x", "y", 50),
    )

    val data2 = Seq(
      Owner(11, "A", 666, 2017, "x", "y", 50),
      Owner(11, "A", 222, 2018, "x", "y", 20),
      Owner(33, "C", 444, 2018, "x", "y", 20),
      Owner(33, "C", 555, 2018, "x", "y", 55),
      Owner(22, "B", 555, 2018, "x", "y", 20),
      Owner(99, "D", 888, 2018, "x", "y", 100),
      Owner(11, "A", 888, 2018, "x", "y", 100),
      Owner(11, "A", 666, 2018, "x", "y", 80),
      Owner(33, "C", 666, 2018, "x", "y", 80),
      Owner(11, "A", 444, 2018, "x", "y", 50),
    )

    val expected = Seq(
      Result(11, "A", 666, 2017, "x", "y", 50, isDiff = false),
      Result(11, "A", 222, 2018, "x", "y", 20, isDiff = false),
      Result(33, "C", 444, 2018, "x", "y", 20, isDiff = false),
      Result(33, "C", 555, 2018, "x", "y", 55, isDiff = true),
      Result(22, "B", 555, 2018, "x", "y", 20, isDiff = false),
      Result(99, "D", 888, 2018, "x", "y", 100, isDiff = false),
      Result(11, "A", 888, 2018, "x", "y", 100, isDiff = false),
      Result(11, "A", 666, 2018, "x", "y", 80, isDiff = false),
      Result(33, "C", 666, 2018, "x", "y", 80, isDiff = false),
      Result(11, "A", 444, 2018, "x", "y", 50, isDiff = false),
    )


    val df1 = spark
      .createDataset(data1)
      .as[Owner]
      .cache()

    val df2 = spark
      .createDataset(data2)
      .as[Owner]
      .cache()
  }

}

最有效的方法是什么?

【问题讨论】:

  • 这能回答你的问题吗? Compare two Spark dataframes
  • 不是真的,这个问题没有考虑隔离记录不同的标志。
  • 我认为下面的 2 个答案效率不高,因为 joinintersect 为所有记录和所有分区创建哈希表并进行比较。至少您可以尝试最简单的解决方案:df1.rdd.zip(df2.rdd).map {case (x,y) => (x, x != y)} 并比较真实数据集的速度。 PS:将单字符字符串替换为char是个好主意,因为char比较非常快
  • @MikhailIonkin 您能否提供完整的答案,以便我接受您的答案。你在这里说得很好!
  • 完成了。我没有真正的数据集,所以我不能证明我的答案更快,但我认为这是根据对小数据集的测试和文档

标签: sql scala apache-spark dataset


【解决方案1】:

我认为这段代码可以帮助您找到答案:

val intersectDF=df1.intersect(df2)
val unionDF=df1.union(df2).dropDuplicates()
val diffDF= unionDF.except(intersectDF)

val intersectDF2=intersectDF.withColumn("isDiff",functions.lit(false))
val diffDF2=diffDF.withColumn("isDiff",functions.lit(true))
val answer=intersectDF2.union(diffDF2)

//Common data between two DataFrame
intersectDF2.show()
//Difference data between two DataFrame
diffDF2.show()
//Your answer
answer.show()

【讨论】:

    【解决方案2】:

    也许这有帮助 -

    进行左连接并将不匹配的列识别为假

     val df1_hash = df1.withColumn("x", lit(0))
        df2.join(df1_hash, df2.columns, "left")
          .select(when(col("x").isNull, false).otherwise(true).as("isDiff") +: df2.columns.map(df2(_)): _*)
          .show(false)
    
        /**
          * +------+---+---+---+----+---+---+------+
          * |isDiff|a  |b  |c  |d   |e  |f  |o_qtty|
          * +------+---+---+---+----+---+---+------+
          * |true  |11 |A  |666|2017|x  |y  |50.0  |
          * |true  |11 |A  |222|2018|x  |y  |20.0  |
          * |true  |33 |C  |444|2018|x  |y  |20.0  |
          * |false |33 |C  |555|2018|x  |y  |55.0  |
          * |true  |22 |B  |555|2018|x  |y  |20.0  |
          * |true  |99 |D  |888|2018|x  |y  |100.0 |
          * |true  |11 |A  |888|2018|x  |y  |100.0 |
          * |true  |11 |A  |666|2018|x  |y  |80.0  |
          * |true  |33 |C  |666|2018|x  |y  |80.0  |
          * |true  |11 |A  |444|2018|x  |y  |50.0  |
          * +------+---+---+---+----+---+---+------+
          */
    

    【讨论】:

      【解决方案3】:

      我认为2其他答案是不高效的,因为joinintersect创建所有记录和所有分区的哈希表,并比较它的全部。至少您可以尝试最简单的解决方案:

        df1.rdd.zip(df2.rdd).map {
          case (x,y) => (x, x != y)
        }
      

      并比较真实数据集的速度。

      另外,最好将单字符字符串替换为 char,因为 char 比较非常快。

      我没有真实数据集,所以我不能证明我的答案更快,但我认为它是根据小型数据集的测试,并根据joinintersectintersectintersectintersectintersectintersect 987654326 @不交换分区与@ 987654327或@或intersect。对不起我的英文

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2022-12-01
        • 1970-01-01
        • 2015-01-04
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2023-04-07
        相关资源
        最近更新 更多