【问题标题】:How to add conditional join in spark如何在火花中添加条件连接
【发布时间】:2018-02-22 09:01:41
【问题描述】:

我有一个类似的数据框连接条件

 

df1.as("main_data") .join(df2.as("mcp"),df1.col("id").equalTo(df2.col("id")) 和 df1.col("name").equalTo(df2.col("name "))

在这个连接中,第二个检查是有条件的

ie df1.col("name").equalTo(df2.col("name")) 应该只执行 如果 includeNameFlag 为假

如何将其添加到我的数据框连接

尝试将条件包含为字符串并与连接一起附加

var joinVar = ""

if(includeNameFlag == false){

    joinVar = """and df1.col("name").equalTo(df2.col("name"))"""

}else{
    joinVar = ""
}

df1.as("main_data")

.join(df2.as("mcp"),df1.col("id").equalTo(df2.col("id"))+ joinVar)

但这没有帮助。它遇到了无法解析 id= id + name =name 之类的错误

尝试使用 when 和 where 条件,但都需要列类型

在数据框连接中使用此条件的任何其他解决方案?

solution

【问题讨论】:

  • 是否有异常或其他输出?
  • 如果我理解正确,您会遇到语法问题(“两者都需要列类型”)。从这个问题中不清楚究竟什么是行不通的。但是,在我看来,您有两个问题: 1. 您试图将 Column 类型与字符串连接起来。 2. and 运算符使用不当。请参阅spark.apache.org/docs/latest/api/java/org/apache/spark/sql/…,您可以使用“condition1.and(condition2)”或“condition1 && condition2”。 “条件 1 和条件 2”无效。
  • 使用spark.sql并根据您的if - else制定查询
  • 数据框中还有其他方式吗?

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


【解决方案1】:

使用DataFrame API 可以轻松完成。这是一个简单的例子:

val df1 = Seq((1, "foo"), (2, "bar")).toDF("id", "name")
val df2 = Seq((1, "bar"), (2, "bar")).toDF("id", "name")

我正在使用类似于equalTo 的等值连接。

一方面,根据你的描述:

val includeNameFlag: Boolean = false
val exprs = (if (!includeNameFlag) Seq("id","name") else Seq("id"))

df1.join(df2, exprs).show
// +---+----+
// | id|name|
// +---+----+
// |  2| bar|
// +---+----+

另一方面:

val includeNameFlag: Boolean = true
val exprs = (if (!includeNameFlag) Seq("id","name") else Seq("id"))

df1.join(df2, exprs).show
// +---+----+----+
// | id|name|name|
// +---+----+----+
// |  1| foo| bar|
// |  2| bar| bar|
// +---+----+----+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-05-31
    • 1970-01-01
    • 2016-12-09
    • 1970-01-01
    • 2017-07-05
    • 2019-02-07
    • 2016-06-05
    • 1970-01-01
    相关资源
    最近更新 更多