Spark中row_number指定全窗口范围失效问题及实现咨询
问题解答:Spark中实现目标row_number逻辑及语句合法性分析
一、先明确:你的原查询语句不符合标准SQL规范
row_number()这类排名窗口函数(包括rank()、dense_rank())的设计逻辑是基于整个分区内的排序生成行号/排名,它们不需要也不支持指定窗口帧(rows between ...)。
Hive之所以能运行,是因为它的窗口函数实现做了宽松兼容,自动忽略了这个多余的窗口范围定义;但Spark SQL(包括1.6和2.0版本)严格遵循ANSI SQL规范,会直接拒绝这种不符合语法规则的写法,这就是你报错的核心原因。
二、Spark中实现你想要的逻辑
其实你原语句的核心需求应该是:在指定分区内,按某列排序后生成全局的行号——这个需求在Spark里根本不需要加窗口范围,直接简化写法即可:
假设你的原Hive语句是:
SELECT your_column1, your_column2, row_number() over (partition by partition_col order by sort_col rows between unbounded preceding and unbounded following) as row_num FROM your_table
在Spark SQL里直接改成:
SELECT your_column1, your_column2, row_number() over (partition by partition_col order by sort_col) as row_num FROM your_table
如果是用Spark Scala API实现,代码示例如下:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.row_number // 定义窗口规则:按partition_col分区,sort_col排序 val windowSpec = Window.partitionBy("partition_col").orderBy("sort_col") // 给DataFrame添加行号列 val resultDF = yourDF.withColumn("row_num", row_number().over(windowSpec))
这样写完全等价于你原Hive语句的逻辑,因为row_number()默认就是对整个分区生效的,窗口范围的指定在这里完全是冗余的。
内容的提问来源于stack exchange,提问作者Siva kumar
相关产品推荐
相关产品推荐

