【问题标题】:UDF in Spark SQL DSLSpark SQL DSL 中的 UDF
【发布时间】:2016-06-21 06:25:50
【问题描述】:

我试图在 Spark SQL 作业中使用 DSL 而不是纯 SQL,但我的 UDF 无法正常工作。

sqlContext.udf.register("subdate",(dateTime: Long)=>dateTime.toString.dropRight(6))

这行不通

rdd1.toDF.join(rdd2.toDF).where("subdate(rdd1(date_time)) === subdate(rdd2(dateTime))")

我还想在这个工作纯 SQL 中添加另一个连接条件

val results=sqlContext.sql("select * from rdd1 join rdd2 on rdd1.id=rdd2.idand subdate(rdd1.date_time)=subdate(rdd2.dateTime)")

感谢您的帮助

【问题讨论】:

    标签: sql apache-spark apache-spark-sql user-defined-functions dsl


    【解决方案1】:

    您传递给where 方法的SQL 表达式不正确,至少有几个原因:

    • ===Column 方法,不是有效的 SQL 相等性。您应该使用单个等号 =
    • 括号表示法 (table(column)) 不是在 SQL 中引用列的有效方法。在这种情况下,它将被识别为函数调用。 SQL 使用点表示法 (table.column)
    • 即使 rdd1rdd2 都不是有效的表别名

    由于看起来列名是明确的,您可以简单地使用以下代码:

    df1.join(df2).where("subdate(date_time) = subdate(dateTime)")
    

    如果不是这样,如果不先提供别名,使用点语法将无法工作。参见例如Usage of spark DataFrame "as" method

    此外,当您一直使用原始 SQL 时,注册 UDF 最有意义。如果要使用DataFrame API,最好直接使用UDF:

    import org.apache.spark.sql.functions.udf
    
    val subdate = udf((dateTime: Long) => dateTime.toString.dropRight(6)) 
    
    val df1 = rdd1.toDF
    val df2 = rdd2.toDF
    
    df1.join(df2, subdate($"date_time") === subdate($"dateTime"))
    

    或者如果列名不明确:

    df1.join(df2, subdate(df1("date_time")) === subdate(df2("date_time")))
    

    最后,对于像这样的简单函数,编写内置表达式比创建 UDF 更好。

    【讨论】:

    • 非常感谢。编写内置表达式是什么意思?使用 sql.Column 包中的类似“substr”的函数?
    • 或多或少。那里有一些微妙之处(并非每个函数都是使用表达式实现的),但不要纠缠于此。如果这有帮助,请不要感谢 - 只需接受和/或支持 :)
    猜你喜欢
    • 2017-09-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-03
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多