【问题标题】:3 LEFT-JOIN in Spark SQL with API JAVA3 使用 API JAVA 在 Spark SQL 中的 LEFT-JOIN
【发布时间】:2021-03-14 09:05:09
【问题描述】:

我有来自 3 个表的 3 个数据集:

Dataset<TABLE1> bbdd_one = map.get("TABLE1").as(Encoders.bean(TABLE1.class)).alias("TABLE1");
Dataset<TABLE2> bbdd_two = map.get("TABLE2").as(Encoders.bean(TABLE2.class)).alias("TABLE2");
Dataset<TABLE3> bbdd_three = map.get("TABLE3").as(Encoders.bean(TABLE3.class)).alias("TABLE3");

我想对其进行三重左连接并将其写入输出 .parquet

sql JOIN 语句类似这样:

SELECT one.field, ........, two.field ....., three.field, ... four.field
FROM TABLE1 one
LEFT JOIN TABLE2 two ON two.field = one.field
LEFT JOIN TABLE3 three ON three.field = one.field AND three.field = one.field
LEFT JOIN TABLE3 four ON four.field = one.field AND four.field = one.otherfield
WHERE one.field = 'whatever'

如何使用 JAVA API 做到这一点?可能吗?我做了一个只有一个连接但有 3 个连接的示例。

PS:我的另一个加入 JAVA API 是:

Dataset<TJOINED> ds_joined = ds_table1
                        .join(ds_table2,
                                JavaConversions.asScalaBuffer(Arrays.asList("fieldInCommon1", "fieldInCommon2", "fieldInCommon3", "fieldInCommon4"))
                                        .seq(),
                                "inner")
                        .select("a lot of fields", ... "more fields")                                                               
                        .as(Encoders.bean(TJOINED.class));

谢谢!

【问题讨论】:

    标签: java apache-spark join apache-spark-sql left-join


    【解决方案1】:

    您是否尝试过链接连接语句? 我不经常用 Java 编写代码,所以这只是一个猜测

    Dataset<TJOINED> ds_joined = ds_table1
        .join(
            ds_table2,
            JavaConversions.asScalaBuffer(Arrays.asList(...)).seq(),
            "left"
        )
        .join(
            ds_table3,
            JavaConversions.asScalaBuffer(Arrays.asList(...)).seq(),
            "left"
        )
        .join(
            ds_table4,
            JavaConversions.asScalaBuffer(Arrays.asList(...)).seq(),
            "left"
        )
        .select(...)
        .as(Encoders.bean(TJOINED.class))
    

    更新:如果我的理解是正确的,ds_table3ds_table4 是相同的,它们是在不同的领域加入的。那么也许这个更新的答案,它是在 Scala 中给出的,因为它是我习惯使用的,可能会实现你想要的。这是完整的工作示例:

    import spark.implicits._
    
    case class TABLE1(f1: Int, f2: Int, f3: Int, f4: Int, f5:Int)
    case class TABLE2(f1: Int, f2: Int, vTable2: Int)
    case class TABLE3(f3: Int, f4: Int, vTable3: Int)
    
    val one = spark.createDataset[TABLE1](Seq(TABLE1(1,2,3,4,5), TABLE1(1,3,4,5,6)))
    //one.show()
    //+---+---+---+---+---+
    //| f1| f2| f3| f4| f5|
    //+---+---+---+---+---+
    //|  1|  2|  3|  4|  5|
    //|  1|  3|  4|  5|  6|
    //+---+---+---+---+---+
    
    val two = spark.createDataset[TABLE2](Seq(TABLE2(1,2,20)))
    //two.show()
    //+---+---+-------+
    //| f1| f2|vTable2|
    //+---+---+-------+
    //|  1|  2|     20|
    //+---+---+-------+
    
    val three = spark.createDataset[TABLE3](Seq(TABLE3(3,4,20), TABLE3(3,5,50)))
    //three.show()
    //+---+---+-------+
    //| f3| f4|vTable3|
    //+---+---+-------+
    //|  3|  4|     20|
    //|  3|  5|     50|
    //+---+---+-------+
    
    val result = one
    .join(two, Seq("f1", "f2"), "left")
    .join(three, Seq("f3", "f4"), "left")
    .join(
      three.withColumnRenamed("f4", "f5").withColumnRenamed("vTable3", "vTable4"),
      Seq("f3", "f5"),
      "left"
    )
    //result.show()
    //+---+---+---+---+---+-------+-------+-------+
    //| f3| f5| f4| f1| f2|vTable2|vTable3|vTable4|
    //+---+---+---+---+---+-------+-------+-------+
    //|  3|  5|  4|  1|  2|     20|     20|     50|
    //|  4|  6|  5|  1|  3|   null|   null|   null|
    //+---+---+---+---+---+-------+-------+-------+
    

    【讨论】:

    • 您的方法很好,但问题是第 2 次和第 3 次 JOIN(同一个表),特别是:LEFT JOIN TABLE3 three ON three.field = one.field AND three.field = one.field (名称匹配)LEFT JOIN TABLE3 4 ONfour.field = one.field ANDfour.field = one.otherfield(主表名称与TABLE3名称不匹配)
    • 抛出以下异常:UNNAMED with exception org.apache.spark.sql.AnalysisException: USING column one.otherfield cannot be resolve on the right side of the join
    • 我假设我们会有这样的东西: TABLE1 (f1, f2, f3, f4, f5) TABLE2 (f1, f2) TABLE3 (f3,f4) 然后在第三次加入你想要执行以下操作:LEFT JOIN TABLE3 four on four.f3 = one.f3 and four.f4 = one.f5 我认为可以通过临时保存前 2 个连接的结果并执行最后一个连接来实现。我习惯于使用 Dataset,所以我并不真正了解这些类型的 Dataset 是如何工作的
    • 您的答案有效!抱歉迟到了,但非常感谢!!
    猜你喜欢
    • 1970-01-01
    • 2010-09-29
    • 2021-10-21
    • 1970-01-01
    • 2014-08-10
    • 2011-06-24
    • 1970-01-01
    • 2016-12-16
    相关资源
    最近更新 更多