【问题标题】:How to join data frame in Apache spark如何在 Apache Spark 中加入数据框
【发布时间】:2016-05-19 15:50:19
【问题描述】:

我有以下两个数据框:

df1

uid   text   frequency
11    a      1
12    a      2
12    b      1

df2

text
a
b
c
d

我想创建一个类似这样的数据框:

输出df

uid  text  frequency
11   a     1
11   b     0
11   c     0
11   d     0
12   a     2
12   b     1
12   c     0
12   d     0

我一直在使用 spark-sql 来编写这样的连接:

 sqlContext.sql("Select uid,df2.text,frequency from df1  right outer join df2 on df1.text= df2.text") 

没有返回正确的结果。

有什么建议吗?

【问题讨论】:

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


    【解决方案1】:

    你必须这样做

    // Find unique combinations of uid and text
    df1.select("uid").distinct.join(df2.distinct)  
      // Left join with df1
      .join(df1, Seq("uid", "text"), "leftouter")
      // Replace missing values with 0
      .withColumn("frequency", coalesce($"frequency", lit(0)))
    

    大致相当于下面的SQL:

    WITH tmp AS (SELECT DISTINCT df1.uid, df2.text FROM df1  JOIN df2)
    SELECT tmp.uid, tmp.text, COALESCE(df1.frequency, 0) AS frequency
    FROM tmp LEFT OUTER JOIN df1
    ON tmp.uid = df1.uid AND tmp.text = df1.text
    

    【讨论】:

    • 从性能的角度来看,与 SQL 查询相比,第一个答案需要花费大量时间。有什么办法可以改善这一点,因为我使用的是数据框而不是表格。
    • SQL查询什么?如果您的意思是 Spark SQL 查询,这里没有区别。两者都应评估为相同的执行计划。
    猜你喜欢
    • 1970-01-01
    • 2018-02-21
    • 1970-01-01
    • 2020-03-18
    • 1970-01-01
    • 1970-01-01
    • 2017-02-13
    • 1970-01-01
    • 2015-06-11
    相关资源
    最近更新 更多