【发布时间】: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