【问题标题】:How to do a self join in Spark 2.3.0? What is the correct syntax?如何在 Spark 2.3.0 中进行自我加入?什么是正确的语法?
【发布时间】:2018-02-21 07:18:00
【问题描述】:

我有以下代码

import org.apache.spark.sql.streaming.Trigger 

val jdf = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "join_test").option("startingOffsets", "earliest").load();   
jdf.createOrReplaceTempView("table")
val resultdf = spark.sql("select * from table as x inner join table as y on x.offset=y.offset")
resultdf.writeStream.outputMode("append").format("console").option("truncate", false).trigger(Trigger.ProcessingTime(1000)).start()

我得到以下异常

org.apache.spark.sql.AnalysisException: cannot resolve '`x.offset`' given input columns: [x.value, x.offset, x.key, x.timestampType, x.topic, x.timestamp, x.partition]; line 1 pos 50;
'Project [*]
+- 'Join Inner, ('x.offset = 'y.offset)
   :- SubqueryAlias x
   :  +- SubqueryAlias table
   :     +- StreamingRelation DataSource(org.apache.spark.sql.SparkSession@15f3f9cf,kafka,List(),None,List(),None,Map(startingOffsets -> earliest, subscribe -> join_test, kafka.bootstrap.servers -> localhost:9092),None), kafka, [key#28, value#29, topic#30, partition#31, offset#32L, timestamp#33, timestampType#34]
   +- SubqueryAlias y
      +- SubqueryAlias table
         +- StreamingRelation DataSource(org.apache.spark.sql.SparkSession@15f3f9cf,kafka,List(),None,List(),None,Map(startingOffsets -> earliest, subscribe -> join_test, kafka.bootstrap.servers -> localhost:9092),None), kafka, [key#28, value#29, topic#30, partition#31, offset#32L, timestamp#33, timestampType#34]

我已经把代码改成了这个

import org.apache.spark.sql.streaming.Trigger 

val jdf = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "join_test").option("startingOffsets", "earliest").load();
val jdf1 = spark.readStream.format("kafka").option("kafka.bootstrap.servers", "localhost:9092").option("subscribe", "join_test").option("startingOffsets", "earliest").load();

jdf.createOrReplaceTempView("table")
jdf1.createOrReplaceTempView("table1")

val resultdf = spark.sql("select * from table inner join table1 on table.offset=table1.offset")

resultdf.writeStream.outputMode("append").format("console").option("truncate", false).trigger(Trigger.ProcessingTime(1000)).start()

这行得通。但是,我不相信这是我正在寻找的解决方案。我希望能够使用原始 SQL 进行自联接,而不是像上面的代码那样制作数据帧的额外副本。那么还有其他方法吗?

【问题讨论】:

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


    【解决方案1】:

    这是一个已知问题,将在 2.4.0 中修复。见https://issues.apache.org/jira/browse/SPARK-23406。现在你可以避免加入相同的 DataFrame 对象。

    【讨论】:

    • 啊!!知道了!感谢那。我很难避免加入相同的 DataFrame 对象,主要是因为我从我们的用户那里获得了原始 sql,并且原始 sql 可以包含任意数量的自连接,所以我必须先解析原始 sql,然后尝试创建任意数量的数据框对象等等,所以它会变成一个有点复杂的东西。我看到已经创建了一个 PR,所以我可以将其挑选到 2.3 中,还是可以编译并使用主分支?这个功能对我们来说非常重要,Spark 2.4 可能还需要 5 个月的时间......所以想知道我是否有运气?
    • 如果您自己构建 Spark,那么您绝对可以将其挑选到您自己的 2.3 分支中。社区没有将其合并到 2.3 的原因主要是此补丁中的更改并非微不足道,并且 Spark 2.3.0 正在投票中。
    • 非常感谢!关于这个的最后一个问题..我应该从 master 或 branch-2.3 创建一个分支吗?我只想要 2.3 加上一个修复自联接的提交。
    • 我觉得从 master 分支在我们的生产中部署有点冒险,因为 spark master 分支将不断发展,我可以运行它不兼容的 jars 问题,例如 kafka 等。所以需要从一个可以像 2.3 一样保持稳定的分支分支出来。
    【解决方案2】:

    您可以使用 DataFrame API join 函数而不是使用 SQL 语法:

    jdf.as("df1").join(jdf.as("df2"), $"df1.offset" === $"df2.offset", "inner")
    

    【讨论】:

    • 完全相同的异常org.apache.spark.sql.AnalysisException: cannot resolve 'df1.offset' given input columns: [df1.partition, df1.timestampType, df1.topic, df1.timestamp, df1.value, df1.offset, df1.key];;
    • 仅供参考,我正在尝试进行流式连接(从我的代码中可以看到)。流式传输还是非流式传输似乎并不重要。自联接的语法似乎在这两种情况下都不起作用。我得到了同样的例外。我也在 Spark 2.2.0 中尝试了非流式版本的自我加入,但我得到了同样的错误!!
    • @user1870400:您在使用流式数据集吗?根据文档,尚不支持两个流数据集之间的连接。
    • 我使用的是 Spark 2.3,它们在 2.3 RC4 中受支持。 Spark 2.3 尚未发布,但可在此处获得候选版本 4 dist.apache.org/repos/dist/dev/spark/v2.3.0-rc4-bin
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-01-20
    • 2016-07-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-04-19
    相关资源
    最近更新 更多