Spark WindowSpec lag函数计算累计分数报错问题排查
解决Spark计算用户累计运行分数的报错问题
首先,咱们来拆解你遇到的问题和背后的原因:
错误根源分析
你碰到的AnalysisException窗口帧不匹配问题,主要是两个原因导致的:
- 窗口帧规则冲突:
lag函数默认只需要访问前一行数据,它要求的窗口帧是RowFrame, -1, -1;但你定义的wSpec2窗口用了rowsBetween(Long.MinValue, 0)(也就是从分区开头到当前行的所有数据),两者的窗口范围规则不兼容,直接触发了报错。 - 自引用未创建的列:你在
withColumn里尝试引用正在创建的CumScore列,这在Spark SQL里是不允许的——这个列还没被注册到DataFrame的元数据中,Spark根本找不到它。
另外,从你的需求来看,计算累计运行分数完全不需要用lag递归,Spark的窗口聚合函数sum就能更简洁高效地搞定这件事。
正确解决方案
结合你“每个用户按日期累计,前一日累计分数加当日分数”的需求,再加上你代码里的乘2逻辑,咱们可以这样修改:
第一步:定义正确的窗口
先按name分区(保证每个用户独立计算),再按date排序(确保日期顺序正确),窗口范围设为从分区起始到当前行:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 先重新创建原始数据 val genre = sc.parallelize(List( ("Alice", "2016-05-01", "action",0), ("Alice", "2016-05-02", "0",1), ("Alice", "2016-05-03", "comedy",0), ("Alice", "2016-05-04", "action",1), ("Alice", "2016-05-05", "action",0), ("Alice", "2016-05-06", "horror",1), ("Bob", "2016-05-01", "art",0), ("Bob", "2016-05-02", "0",1), ("Bob", "2016-05-03", "0",0), ("Bob", "2016-05-04", "art",0), ("Bob", "2016-05-05", "comedy",1), ("Bob", "2016-05-06", "action",0) )).toDF("name", "date", "genre","score") // 定义正确的窗口规则 val wSpec = Window.partitionBy("name").orderBy("date").rowsBetween(Long.MinValue, 0)
第二步:计算累计分数
用sum窗口函数直接计算累计值,同时保留你代码里的乘2逻辑:
val result = genre.withColumn( "CumScore", sum(score * 2).over(wSpec) ) result.show()
最终输出结果
执行后会得到每个用户按日期累计的分数:
+-----+----------+------+-----+--------+ | name| date| genre|score|CumScore| +-----+----------+------+-----+--------+ |Alice|2016-05-01|action| 0| 0| |Alice|2016-05-02| 0| 1| 2| |Alice|2016-05-03|comedy| 0| 2| |Alice|2016-05-04|action| 1| 4| |Alice|2016-05-05|action| 0| 4| |Alice|2016-05-06|horror| 1| 6| | Bob|2016-05-01| art| 0| 0| | Bob|2016-05-02| 0| 1| 2| | Bob|2016-05-03| 0| 0| 2| | Bob|2016-05-04| art| 0| 2| | Bob|2016-05-05|comedy| 1| 4| | Bob|2016-05-06|action| 0| 4| +-----+----------+------+-----+--------+
补充说明
如果之后你有更复杂的累计逻辑(不是简单求和),可以考虑用Spark的递归CTE来实现,但对于这种基础的累计求和,sum窗口函数是性能最优、代码最简洁的选择。
内容的提问来源于stack exchange,提问作者Masterbuilder
相关产品推荐
相关产品推荐

