Spark Scala中如何用collect_list按时间戳排序生成有序列表?
Spark分组后collect_list无法按时间排序的问题及解决办法
问题原因
- 分布式计算特性导致全局排序失效:Spark是分布式计算框架,全局排序后的数据会重新分区,分组操作在各个分区独立执行,同组的数据可能分散在不同分区,且分区内的同组元素无法保证是按
creation_date有序的,collect_list只会按数据在分区内的存储顺序收集元素。 collect_list默认无排序逻辑:这个聚合函数本身不对输入元素做排序处理,只是单纯将同组元素聚合为列表,所以即使提前全局排序,也无法保证分组后的列表有序。
解决办法
方法1:Spark 3.0+ 直接使用带排序的collect_list
Spark 3.0及以上版本支持在collect_list中指定排序规则,直接在聚合阶段按creation_date降序排列元素:
db.groupBy("ids") .agg(collect_list("names").orderBy(desc("creation_date")).alias("items")) .select("ids", "items")
方法2:低版本Spark 用struct+sort_array实现
先将creation_date和names封装成结构体,收集结构体列表后,按creation_date降序排序结构体,最后提取names字段:
db.groupBy("ids") .agg(collect_list(struct("creation_date", "names")).alias("temp_list")) .select( "ids", sort_array($"temp_list", asc = false).names.alias("items") )
为什么之前的尝试无效
你之前直接对全量数据按creation_date排序后再分组,但Spark的分组操作依赖分区,排序后的数据会被重新分配到不同分区,每个分区内的同组数据顺序无法保证,collect_list只是按分区内的现有顺序聚合元素,因此无法得到按时间从新到旧排列的列表。
内容的提问来源于stack exchange,提问作者Jericho Sims
相关产品推荐
相关产品推荐

