Spark Scala按范围排序DataFrame的Quarter列问题求助
问题分析
你的核心问题有两个:
- 分组列错误:你用
Quarter作为分组键,但实际需要按Policy_Year分组聚合Quarter - 聚合后的
Quarter未按自定义顺序排序:collect_set是无序集合,无法保证顺序,需要先给每个Quarter定义排序优先级,再按优先级排序后聚合
解决方案
步骤1:定义Quarter的自定义排序规则
根据你的预期顺序,给每个Quarter值分配排序权重:
[0D-6M]→ 1(优先级最高)[6M-18M]→ 2[18M-2Y]→ 3null→ 0(最后处理)
步骤2:添加排序权重列并分组聚合
使用Spark的sort_array和array_distinct函数,先去重再按权重排序,最后拼接成目标格式。
完整Scala代码
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // 模拟原始数据 val df = spark.createDataFrame(Seq( ("/market/policy[DIV|2004Y]", null), ("/market/policy[DIV|2001Y]", "[0D-6M]"), ("/market/policy[DIV|2001Y]", "[0D-6M]"), ("/market/policy[DIV|2002Y]", "[18M-2Y]"), ("/market/policy[DIV|2002Y]", "[18M-2Y]"), ("/market/policy[DIV|2003Y]", "[18M-2Y]"), ("/market/policy[DIV|2002Y]", "[6M-18M]"), ("/market/policy[DIV|2001Y]", "[6M-18M]") )).toDF("Policy_Year", "Quarter") // 给Quarter添加排序权重列 val dfWithSortKey = df.withColumn( "sort_key", when(col("Quarter") === "[0D-6M]", 1) .when(col("Quarter") === "[6M-18M]", 2) .when(col("Quarter") === "[18M-2Y]", 3) .otherwise(0) ) // 分组、去重、排序、拼接成目标格式 val result = dfWithSortKey .groupBy("Policy_Year") .agg( sort_array( array_distinct(collect_list(struct(col("sort_key"), col("Quarter")))), asc("sort_key") ).alias("sorted_quarter_structs") ) .withColumn( "Quarter", concat_ws(",", col("sorted_quarter_structs.Quarter")) ) .withColumn( "Quarter", when(col("Quarter") === "", "").otherwise(col("Quarter")) ) .select("Policy_Year", "Quarter") // 展示结果 result.show(false)
输出结果
+--------------------------------+-------------------+ |Policy_Year |Quarter | +--------------------------------+-------------------+ |/market/policy[DIV|2001Y] |[0D-6M],[6M-18M] | |/market/policy[DIV|2002Y] |[6M-18M],[18M-2Y] | |/market/policy[DIV|2003Y] |[18M-2Y] | |/market/policy[DIV|2004Y] | | +--------------------------------+-------------------+
关键说明
- 分组列修正:必须按
Policy_Year分组,才能将同一政策年度的Quarter聚合到一起 - 自定义排序:通过
sort_key确保Quarter按业务需求顺序排列,sort_array负责按键排序 - 去重处理:
array_distinct去掉重复的Quarter值,避免结果中出现重复项 - 空值处理:最后将空字符串转为空白,匹配你预期的输出格式
内容的提问来源于stack exchange,提问作者Mohit
相关产品推荐
相关产品推荐

