【问题标题】:Spark SQL: Unable to use aggregate within a window functionSpark SQL:无法在窗口函数中使用聚合
【发布时间】:2017-11-30 19:22:35
【问题描述】:

我使用此 SQL 为数据集创建 session_id。如果用户处于非活动状态超过 30 分钟(30*60 秒),则会分配一个新的 session_id 我是 Spark SQL 的新手,并尝试使用 Spark SQL 上下文复制相同的过程。但是我遇到了一些错误。

session_id 遵循命名约定:
用户 ID_1,
用户 ID_2,
userid_3,...

SQL(日期以秒为单位):

CREATE TABLE tablename_with_session_id AS 
    SELECT * , userid || '_' || SUM(new_session) OVER (PARTITION BY userid ORDER BY date asc, new_session desc rows unbounded preceding) AS session_id 
        FROM 
                (SELECT *, 
                    CASE 
                    WHEN (date - LAG(date) OVER (PARTITION BY userid ORDER BY date) >= 30 * 60) 
                    THEN 1 
                    WHEN row_number() over (partition by userid order by date) = 1 
                    THEN 1 
                    ELSE 0 
                    END as new_session 
                    FROM 
                        tablename 
                ) 
    order by date;

我尝试在 Spark-Scala 中使用相同的 SQL:

    val sqlContext = new org.apache.spark.sql.SQLContext(sc) 

    val tableSessionID = sqlContext.sql("SELECT * , CONCAT(userid,'_',SUM(new_session)) OVER (PARTITION BY userid ORDER BY date asc, new_session desc rows unbounded preceding) AS new_session_id FROM
(SELECT *, CASE WHEN (date - LAG(date) OVER (PARTITION BY userid ORDER BY date) >= 30 * 60) THEN 1 WHEN row_number() over (partition by userid order by date) = 1 THEN 1 ELSE 0 END as new_session FROM clickstream) order by date") 

建议在窗口函数中包装 Spark SQL 表达式 ..sum(new_session).. 的一些错误。

我尝试使用多个数据框:

val temp1 = sqlContext.sql("SELECT *, CASE WHEN (date - LAG(date) OVER (PARTITION BY userid ORDER BY date) >= 30 * 60) THEN 1 WHEN row_number() over (partition by userid order by date) = 1 THEN 1 ELSE 0 END as new_session FROM clickstream")

temp1.registerTempTable("clickstream_temp1")

val temp2 = sqlContext.sql("SELECT * , SUM(new_session) OVER (PARTITION BY userid ORDER BY date asc, new_session desc rows unbounded preceding) AS s_id FROM clickstream_temp1")

temp2.registerTempTable("clickstream_temp2")

val temp3 = sqlContext.sql("SELECT * , CONCAT(userid,'_',s_id) OVER (PARTITION BY userid ORDER BY date asc, new_session desc rows unbounded preceding) AS new_session_id FROM clickstream_temp2")

仅对上述语句返回错误。 'val temp3 = ...' CONCAT(userid,'_',s_id) 不能在窗口函数中使用。

解决方法是什么?有其他选择吗?

谢谢

【问题讨论】:

  • CONCAT(userid, '_', SUM(new_session) OVER (PARTITION BY ...))

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


【解决方案1】:

要将 concat 与 spark 窗口函数一起使用,您需要使用用户定义的聚合函数 (UDAF)。窗口函数不能直接使用 concat 函数。

//Extend UserDefinedAggregateFunction to write custom aggregate function
//You can also specify any constructor arguments. For instance you can have
//CustomConcat(arg1: Int, arg2: String)
class CustomConcat() extends org.apache.spark.sql.expressions.UserDefinedAggregateFunction {
import org.apache.spark.sql.types._
import org.apache.spark.sql.expressions.MutableAggregationBuffer
import org.apache.spark.sql.Row
// Input Data Type Schema
def inputSchema: StructType = StructType(Array(StructField("description", StringType)))
// Intermediate Schema
def bufferSchema = StructType(Array(StructField("groupConcat", StringType)))
// Returned Data Type.
def dataType: DataType = StringType
// Self-explaining
def deterministic = true
// This function is called whenever key changes
def initialize(buffer: MutableAggregationBuffer) = {buffer(0) = " ".toString}
// Iterate over each entry of a group
def update(buffer: MutableAggregationBuffer, input: Row) = { buffer(0) = buffer.getString(0) + input.getString(0) }
// Merge two partial aggregates
def merge(buffer1: MutableAggregationBuffer, buffer2: Row) = { buffer1(0) = buffer1.getString(0) + buffer2.getString(0) }
// Called after all the entries are exhausted.
def evaluate(buffer: Row) = {buffer.getString(0)}
}
val newdescription = new CustomConcat
val newdesc1=newdescription($"description").over(windowspec)

您可以使用 newdesc1 作为聚合函数在窗口函数中进行连接。 有关更多信息,您可以查看: databricks udaf 我希望这能回答你的问题。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-27
    • 1970-01-01
    相关资源
    最近更新 更多