Apache Spark窗口函数FIRST_VALUE失效问题求助
问题分析与解决方案
我来帮你理清楚为什么FIRST_VALUE()窗口函数版本没生效,以及怎么正确用它实现你的需求~
为什么你的dataset2没有生效?
你当前的FIRST_VALUE()用法只是给每一行新增了一个new列,值是当前窗口分区内(按ID分组、VALUEE非空优先排序)的第一个符合条件的VALUEE,但并没有对每个ID只保留一行。执行drop("new")后,原数据集的所有行都还在,自然达不到你要的“每个ID唯一记录”的效果。
而row_number()版本能生效,是因为它给每个ID分组内的行按规则编了号,你筛选出编号为1的行,就直接实现了每个ID只留一行的目标。
正确使用FIRST_VALUE()实现需求的方法
我们需要结合去重或筛选逻辑,让FIRST_VALUE()的结果帮我们保留每个ID的唯一记录,这里提供两种可行写法:
方法一:标记目标值后筛选去重
先给每个行标记出对应ID的优先VALUEE,再筛选出符合条件的行并去重:
WindowSpec windowSpec = Window.partitionBy(dataset.col("ID")) .orderBy(dataset.col("VALUEE").asc_nulls_last()); Dataset<Row> dataset2 = dataset .withColumn("target_value", functions.first("VALUEE", true).over(windowSpec)) // 筛选:当前行VALUEE等于目标值,或者目标值为空(该ID全是null)则保留任意一行 .where(functions.col("VALUEE").equalTo(functions.col("target_value")) .or(functions.col("target_value").isNull())) .drop("target_value") .distinct(); // 去重同一ID下重复的符合条件的行 dataset2.show();
方法二:扩展窗口范围后去重
把窗口范围设置为整个ID分区,这样每个ID下的所有行的new列值都相同,再按ID和new去重:
WindowSpec windowSpecFull = Window.partitionBy(dataset.col("ID")) .orderBy(dataset.col("VALUEE").asc_nulls_last()) .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing); Dataset<Row> dataset2 = dataset .withColumn("new", functions.first("VALUEE", true).over(windowSpecFull)) .dropDuplicates("ID", "new") // 同一ID的行new值相同,去重后只剩一行 .drop("new"); dataset2.show();
补充:改进你的GroupBy版本(保留OTHER列)
你之前的dataset0用groupBy后丢失了OTHER列,其实可以给OTHER列也用first聚合(因为同一ID的OTHER值相同):
Dataset<Row> dataset0 = dataset.groupBy("ID") .agg(functions.first("VALUEE", true).alias("VALUEE"), functions.first("OTHER", true).alias("OTHER")); dataset0.show();
总结
row_number()的核心是给每个分区的行编号,取第一行直接实现去重;FIRST_VALUE()是窗口内的聚合函数,本身不会减少行数,必须配合筛选、去重等逻辑才能实现每个ID只留一行的需求。
内容的提问来源于stack exchange,提问作者raphaelauv
相关产品推荐
相关产品推荐

