You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.25 04:24:04