Spark Scala中按列条件获取最新数据失效问题
解决Spark Scala中分组取最新数据后出现多行的问题
嘿,我来帮你搞定这个问题!你遇到的同一分组出多条数据的情况,核心是代码逻辑没真正实现“每组仅留最新一行”的需求,反而保留了组内所有符合过滤条件的行。咱们一步步拆解:
问题出在哪?
你的现有代码有两个关键漏洞:
- 窗口函数没锁定最新行:你用
first("FFAction|!|").over(windowSpec2)只是给组里每一行都贴了个“组内最新行的FFAction”标签,但并没有筛选出仅那一行最新数据。 - 过滤逻辑留了多行:你的过滤条件是判断当前行的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
相关产品推荐
相关产品推荐

