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

Spark Scala中按列条件获取最新数据失效问题

解决Spark Scala中分组取最新数据后出现多行的问题

嘿,我来帮你搞定这个问题!你遇到的同一分组出多条数据的情况,核心是代码逻辑没真正实现“每组仅留最新一行”的需求,反而保留了组内所有符合过滤条件的行。咱们一步步拆解:

问题出在哪?

你的现有代码有两个关键漏洞:

  1. 窗口函数没锁定最新行:你用first("FFAction|!|").over(windowSpec2)只是给组里每一行都贴了个“组内最新行的FFAction”标签,但并没有筛选出仅那一行最新数据。
  2. 过滤逻辑留了多行:你的过滤条件是判断当前行的FFAction是否符合规则,同时参考组内最新行的FFAction——这就导致组里所有满足条件的行都会被留下,比如SourceID=195的分组里,两行都是O|!|,所以都通过了过滤,最终出现重复分组。

正确的实现方式

咱们得先拿到每组的最新一行,再对这些最新数据做FFAction过滤。最稳妥的方法是给每行加行号,用行号锁定最新行:

修改后的代码

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

// 定义窗口:按OrganizationID+SourceID分区,按时间戳降序排序
val windowSpec = Window.partitionBy("OrganizationID", "SourceID")
  .orderBy($"TimeStamp".desc) // 这里直接用TimeStamp字段即可,它已经是标准timestamp格式

// 三步搞定:加行号→留最新行→过滤FFAction
val latestForEachKey = latestForEachKey1
  .withColumn("row_num", row_number().over(windowSpec))
  .filter($"row_num" === 1) // 只保留每组最新的那一行
  .filter($"FFAction|!|" isin ("I|!|", "O|!|", "D|!|")) // 简化你的FFAction过滤条件
  .drop("row_num", "TimeStamp")

代码说明

  • 加行号:row_number()会给每个分组内的行按时间从新到旧分配行号,最新的行号是1。
  • 锁定最新行:filter($"row_num" === 1)确保每个分组只留最顶部的那一行(也就是最新数据)。
  • 简化过滤:用isin替代多个OR,代码更简洁,效果完全一样。

这样处理后,SourceID=195的分组只会保留时间为2018-05-03T09:52:48+00:00的那一行,其他分组也都会只有一条数据,完美符合你的预期!

内容的提问来源于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 06:32:58