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

Spark Scala按范围排序DataFrame的Quarter列问题求助

问题分析

你的核心问题有两个:

  1. 分组列错误:你用Quarter作为分组键,但实际需要按Policy_Year分组聚合Quarter
  2. 聚合后的Quarter未按自定义顺序排序:collect_set是无序集合,无法保证顺序,需要先给每个Quarter定义排序优先级,再按优先级排序后聚合
解决方案

步骤1:定义Quarter的自定义排序规则

根据你的预期顺序,给每个Quarter值分配排序权重:

  • [0D-6M] → 1(优先级最高)
  • [6M-18M] → 2
  • [18M-2Y] → 3
  • null → 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 06:24:52