【问题标题】:Cleanest, most efficient syntax to perform DataFrame self-join in Spark在 Spark 中执行 DataFrame 自联接的最简洁、最有效的语法
【发布时间】:2016-07-14 21:42:13
【问题描述】:

在标准 SQL 中,当您将表连接到自身时,您可以为表创建别名以跟踪您所引用的列:

SELECT a.column_name, b.column_name...
FROM table1 a, table1 b
WHERE a.common_field = b.common_field;

我可以想到两种使用 Spark DataFrame API 实现相同目标的方法:

解决方案 #1:重命名列

对于 this question 的答复有几种不同的方法。这只是用特定后缀重命名所有列:

df.toDF(df.columns.map(_ + "_R"):_*)

例如你可以这样做:

df.join(df.toDF(df.columns.map(_ + "_R"):_*), $"common_field" === $"common_field_R")

解决方案 #2:将引用复制到 DataFrame

另一个简单的解决方案是这样做:

val df: DataFrame = ....
val df_right = df

df.join(df_right, df("common_field") === df_right("common_field"))

这两种解决方案都有效,我可以看到每种解决方案在某些情况下都很有用。两者之间有什么我应该注意的内部差异吗?

【问题讨论】:

    标签: apache-spark dataframe apache-spark-sql


    【解决方案1】:

    至少有两种不同的方法可以通过别名来解决这个问题:

    df.as("df1").join(df.as("df2"), $"df1.foo" === $"df2.foo")
    

    或使用基于名称的相等连接:

    // Note that it will result in ambiguous column names
    // so using aliases here could be a good idea as well.
    // df.as("df1").join(df.as("df2"), Seq("foo"))
    
    df.join(df, Seq("foo"))  
    

    一般来说,列重命名虽然最丑陋,但却是所有版本中最安全的做法。存在一些与列分辨率相关的错误(不久前的we found one on SO),如果您使用原始表达式,解析器之间的一些细节可能会有所不同(HiveContext / 标准SQLContext)。

    我个人更喜欢使用别名,因为它们与惯用的 SQL 相似,并且能够在特定 DataFrame 对象的范围之外使用。

    关于性能,除非您对接近实时的处理感兴趣,否则应该没有任何性能差异。所有这些都应该生成相同的执行计划。

    【讨论】:

    • DataFrame.as 比较新吗?
    • 不,它至少从 1.3 开始就存在了。它只是不那么常用。不要误认为as[U]Datasets 一起使用。 Scala 和 Python 有一个替代的 alias 方法可以实现相同的目标。
    • 奇怪的是,除了索引之外,我在 scaladocs 中没有看到任何对它的引用。我看到了as[U],但即使我回到早期版本,我也没有在DataFrame 中看到as。我在Column 中看到了as,只是在DataFrame 中没有。
    • 检查语言综合查询部分(第8位左右)或来源github.com/apache/spark/blob/…:)
    • 这么多可以对你做到这一点:)
    猜你喜欢
    • 2022-01-03
    • 1970-01-01
    • 2021-09-27
    • 1970-01-01
    • 2015-11-29
    • 1970-01-01
    • 1970-01-01
    • 2010-12-17
    • 2014-11-21
    相关资源
    最近更新 更多