【问题标题】:Left Anti join in Spark?左反加入 Spark?
【发布时间】:2017-08-28 10:52:17
【问题描述】:

我已经定义了两个这样的表:

 val tableName = "table1"
    val tableName2 = "table2"

    val format = new SimpleDateFormat("yyyy-MM-dd")
      val data = List(
        List("mike", 26, true),
        List("susan", 26, false),
        List("john", 33, true)
      )
    val data2 = List(
        List("mike", "grade1", 45, "baseball", new java.sql.Date(format.parse("1957-12-10").getTime)),
        List("john", "grade2", 33, "soccer", new java.sql.Date(format.parse("1978-06-07").getTime)),
        List("john", "grade2", 32, "golf", new java.sql.Date(format.parse("1978-06-07").getTime)),
        List("mike", "grade2", 26, "basketball", new java.sql.Date(format.parse("1978-06-07").getTime)),
        List("lena", "grade2", 23, "baseball", new java.sql.Date(format.parse("1978-06-07").getTime))
      )

      val rdd = sparkContext.parallelize(data).map(Row.fromSeq(_))
      val rdd2 = sparkContext.parallelize(data2).map(Row.fromSeq(_))
      val schema = StructType(Array(
        StructField("name", StringType, true),
        StructField("age", IntegerType, true),
        StructField("isBoy", BooleanType, false)
      ))
    val schema2 = StructType(Array(
        StructField("name", StringType, true),
        StructField("grade", StringType, true),
        StructField("howold", IntegerType, true),
        StructField("hobby", StringType, true),
        StructField("birthday", DateType, false)
      ))

      val df = sqlContext.createDataFrame(rdd, schema)
      val df2 = sqlContext.createDataFrame(rdd2, schema2)
      df.createOrReplaceTempView(tableName)
      df2.createOrReplaceTempView(tableName2)

我正在尝试构建查询以从 table1 返回 table2 中没有匹配行的行。 我试过用这个查询来做:

Select * from table1 LEFT JOIN table2 ON table1.name = table2.name AND table1.age = table2.howold AND table2.name IS NULL AND table2.howold IS NULL

但这只是给了我 table1 中的所有行:

列表({"name":"john","age":33,"isBoy":true}, {"name":"susan","age":26,"i​​sBoy":false}, {"name":"mike","age":26,"i​​sBoy":true})

如何在 Spark 中高效地进行这种类型的连接?

我正在寻找一个 SQL 查询,因为我需要能够指定要在两个表之间进行比较的列,而不是像在其他推荐问题中那样逐行比较。就像使用减法一样,除了等。

【问题讨论】:

标签: scala apache-spark


【解决方案1】:

您可以使用内置函数except (我会使用你提供的代码,但你没有包含导入,所以我不能只是 c/p 它:()

val a = sc.parallelize(Seq((1,"a",123),(2,"b",456))).toDF("col1","col2","col3")
val b= sc.parallelize(Seq((4,"a",432),(2,"t",431),(2,"b",456))).toDF("col1","col2","col3")

scala> a.show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   a| 123|
|   2|   b| 456|
+----+----+----+


scala> b.show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   4|   a| 432|
|   2|   t| 431|
|   2|   b| 456|
+----+----+----+

scala> a.except(b).show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   1|   a| 123|
+----+----+----+

【讨论】:

  • 我正在寻找一个 SQL 查询,因为我需要能够指定要在两个表之间进行比较的列,而不仅仅是逐行比较
【解决方案2】:

您可以使用“左反”连接类型 - 使用 DataFrame API 或 SQL(DataFrame API 支持 SQL 支持的所有内容,包括您需要的任何连接条件):

数据帧 API:

df.as("table1").join(
  df2.as("table2"),
  $"table1.name" === $"table2.name" && $"table1.age" === $"table2.howold",
  "leftanti"
)

SQL:

sqlContext.sql(
  """SELECT table1.* FROM table1
    | LEFT ANTI JOIN table2
    | ON table1.name = table2.name AND table1.age = table2.howold
  """.stripMargin)

注意:还值得注意的是,有一种更短、更简洁的方法来创建示例数据,无需单独指定架构,使用元组和隐式 toDF 方法,然后“修复”需要时自动推断的架构:

import spark.implicits._
val df = List(
  ("mike", 26, true),
  ("susan", 26, false),
  ("john", 33, true)
).toDF("name", "age", "isBoy")

val df2 = List(
  ("mike", "grade1", 45, "baseball", new java.sql.Date(format.parse("1957-12-10").getTime)),
  ("john", "grade2", 33, "soccer", new java.sql.Date(format.parse("1978-06-07").getTime)),
  ("john", "grade2", 32, "golf", new java.sql.Date(format.parse("1978-06-07").getTime)),
  ("mike", "grade2", 26, "basketball", new java.sql.Date(format.parse("1978-06-07").getTime)),
  ("lena", "grade2", 23, "baseball", new java.sql.Date(format.parse("1978-06-07").getTime))
).toDF("name", "grade", "howold", "hobby", "birthday").withColumn("birthday", $"birthday".cast(DateType))

【讨论】:

    【解决方案3】:

    在 SQL 中,您可以简单地查询到下面(不确定它是否在 SPARK 中有效)

    Select * from table1 LEFT JOIN table2 ON table1.name = table2.name AND table1.age = table2.howold where table2.name IS NULL 
    

    这将返回 table1 中连接失败的所有行

    【讨论】:

    • 这行不通。 where 子句在连接操作之前应用,因此不会产生预期的效果。
    • @Hafthor 这将起作用,并且通常优化器会生成相同的左反连接计划。
    【解决方案4】:

    你可以使用左反。

    dfRcc20.as("a").join(dfClientesDuplicados.as("b")
      ,col("a.eteerccdiid")===col("b.eteerccdiid")&&
        col("a.eteerccdinr")===col("b.eteerccdinr")
      ,"left_anti")
    

    【讨论】:

      猜你喜欢
      • 2020-03-14
      • 1970-01-01
      • 1970-01-01
      • 2019-01-02
      • 2011-03-14
      • 2010-10-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多