Spark DataFrame带partitionBy的Window函数使用问题及需求实现
Spark DataFrame 分组取最新记录并合并的解决方案
看起来你已经搞定了第一步的分组取最新逻辑,接下来的第二步其实可以拆成两个清晰的子任务来实现,我来给你梳理下正确的代码思路和实现方式:
核心思路拆解
- 保留第一步中符合条件的记录:第一步处理后的结果里,FFAction为
I|!|的记录已经是按6列分组后的最新数据,直接提取保留即可。 - 重新处理O/D类型的记录:对第一步结果里FFAction为
O|!|或D|!|的记录,按照指定的5列重新分组,再次用Window函数按TimeStamp降序取每组的最新记录。 - 合并结果:把上述两部分数据合并,得到最终目标输出。
具体代码实现(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
相关产品推荐
相关产品推荐

