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

Spark DataFrame带partitionBy的Window函数使用问题及需求实现

Spark DataFrame 分组取最新记录并合并的解决方案

看起来你已经搞定了第一步的分组取最新逻辑,接下来的第二步其实可以拆成两个清晰的子任务来实现,我来给你梳理下正确的代码思路和实现方式:

核心思路拆解

  1. 保留第一步中符合条件的记录:第一步处理后的结果里,FFAction为I|!|的记录已经是按6列分组后的最新数据,直接提取保留即可。
  2. 重新处理O/D类型的记录:对第一步结果里FFAction为O|!|或D|!|的记录,按照指定的5列重新分组,再次用Window函数按TimeStamp降序取每组的最新记录。
  3. 合并结果:把上述两部分数据合并,得到最终目标输出。

具体代码实现(Scala版本)

首先导入必要的Spark函数:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

假设你第一步处理后的结果DataFrame名为first_step_df,执行以下步骤:

// 1. 提取第一步中FFAction为"I|!|"的记录,直接保留
val df_I = first_step_df.filter(col("FFAction|!|") === "I|!|")

// 2. 处理FFAction为"O|!|"或"D|!|"的记录:按5列分组取最新
val od_window = Window.partitionBy(
  "uniqueFundamentalSet", 
  "PeriodId", 
  "SourceId", 
  "StatementTypeCode", 
  "StatementCurrencyId"
).orderBy(col("TimeStamp").desc)

val df_OD = first_step_df
  .filter(col("FFAction|!|").isin("O|!|", "D|!|"))
  .withColumn("row_rank", row_number().over(od_window))
  .filter(col("row_rank") === 1)
  .drop("row_rank")

// 3. 合并两个结果集(使用unionByName避免列顺序问题)
val final_result_df = df_I.unionByName(df_OD)

Python版本代码(如果用PySpark)

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 1. 提取I类型记录
df_I = first_step_df.filter(F.col("FFAction|!|") == "I|!|")

# 2. 处理O/D类型记录
od_window = Window.partitionBy(
    "uniqueFundamentalSet", 
    "PeriodId", 
    "SourceId", 
    "StatementTypeCode", 
    "StatementCurrencyId"
).orderBy(F.col("TimeStamp").desc())

df_OD = first_step_df.filter(F.col("FFAction|!|").isin("O|!|", "D|!|")) \
    .withColumn("row_rank", F.row_number().over(od_window)) \
    .filter(F.col("row_rank") == 1) \
    .drop("row_rank")

# 3. 合并结果
final_result_df = df_I.unionByName(df_OD)

关键注意点

  • 使用unionByName而非union:确保两个DataFrame的列名完全匹配,避免因列顺序不同导致的数据错位问题。
  • 排序稳定性:如果同一分组内存在多个TimeStamp相同的记录,row_number()会随机选取一条;如果需要固定选取规则,可以在orderBy中添加额外的排序字段(比如UpdateReason_updateReasonId)。
  • 性能优化:如果第一步的结果数据量较大,可以考虑先过滤O/D类型再进行Window操作,减少计算量。

内容的提问来源于stack exchange,提问作者Atharv Thakur

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 08:57:48