Spark窗口函数是否按分区独立运行?分组取最新行出现重复如何解决
问题1:窗口函数返回重复行的成因
你关于“重复行数和分区数量一致”的猜测是错误的,该问题和Spark任务的分区数量没有关系,核心原因是窗口函数不会改变返回的总行数,只会为原数据集的每一行追加计算结果。
你定义的窗口规则是按file_date、some_guid分组,按item_time倒序排序,调用first('item_time').over(daily_window)时,Spark会为该分组内的每一行都赋值为该分组排序后的第一个item_time值。如果你的原分组内有16行数据,最终就会返回16行、所有行item_time都相同的结果,和你观察到的现象完全吻合。
如果要实现“每个分组只返回最新1行”的需求,正确写法是先给分组内的行打行号,再过滤行号为1的记录,示例代码如下:
from pyspark.sql import Window from pyspark.sql.functions import row_number, col daily_window = Window.partitionBy('file_date', 'some_guid').orderBy(col('item_time').desc()) df.withColumn('rn', row_number().over(daily_window)) \ .filter(col('rn') == 1) \ .drop('rn') \ .show()
执行后每个file_date、some_guid分组就只会返回1行最新记录。
问题2:带col4字段的需求实现
分两种常见场景处理:
- 场景1:先筛选所有
col4 = data4的记录,再取其中最新的1行
直接过滤后排序取第一条即可:df.filter(col('col4') == 'data4') \ .orderBy(col('item_time').desc()) \ .limit(1) \ .show() - 场景2:取每个
file_date、some_guid分组下的最新行,同时保留该行对应的col4值
直接复用问题1的行号方案即可,过滤后会自动带出最新行对应的col4字段值:daily_window = Window.partitionBy('file_date', 'some_guid').orderBy(col('item_time').desc()) df.withColumn('rn', row_number().over(daily_window)) \ .filter(col('rn') == 1) \ .drop('rn') \ .select('file_date', 'some_guid', 'item_time', 'col4') \ .show()
内容的提问来源于stack exchange,提问作者Vladimir Shadrin
相关产品推荐
相关产品推荐

