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

如何在Scala/Spark中为DataFrame生成全量去重日期序列列?

使用Scala/Spark生成全局唯一序列列

需求说明

现有输入DataFrame inputDF,包含一列存储天数序列的字段days (seq[String]),需要新增一列all days (seq[String]),该列的值为所有行中出现过的唯一天数集合,且每行的该列值保持一致。

输入示例

+---------------------+
|days (seq[String])   |
+---------------------+
|[sat, sun]           |
|[mon, wed]           |
|[fri ]               |
|[fri, sat]           |
|[mon, sun, sat]      |
+---------------------+

输出示例

+---------------------+----------------------------+
|days (seq[String])   |all days (seq[String])      |
+---------------------+----------------------------+
|[sat, sun]           |[sat, sun, mon, wed, fri]   |
|[mon, wed]           |[sat, sun, mon, wed, fri]   |
|[fri]                |[sat, sun, mon, wed, fri]   |
|[fri, sat]           |[sat, sun, mon, wed, fri]   |
|[mon, sun, sat]      |[sat, sun, mon, wed, fri]   |
+---------------------+----------------------------+

实现代码

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

object GlobalDaysApp {
  def main(args: Array[String]): Unit = {
    // 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("GlobalUniqueDays")
      .master("local[*]") // 本地测试环境使用,生产环境移除该行
      .getOrCreate()
    import spark.implicits._

    // 构造输入DataFrame
    val inputDF = Seq(
      Seq("sat", "sun"),
      Seq("mon", "wed"),
      Seq("fri "),
      Seq("fri", "sat"),
      Seq("mon", "sun", "sat")
    ).toDF("days (seq[String])")

    // 1. 提取全局唯一天数:展开序列→去除空格→去重→收集为有序数组
    val allUniqueDays = inputDF
      .select(explode(col("days (seq[String])")).alias("day"))
      .select(trim(col("day")).alias("day"))
      .distinct()
      .collect()
      .map(_.getString(0))
      .sorted // 可选,保证输出顺序和示例一致

    // 2. 广播全局唯一天数集合,避免重复传输
    val broadcastAllDays = spark.sparkContext.broadcast(allUniqueDays)

    // 3. 给原DataFrame添加新列
    val outputDF = inputDF.withColumn(
      "all days (seq[String])",
      lit(broadcastAllDays.value).cast("array<string>")
    )

    // 查看结果
    outputDF.show(false)

    spark.stop()
  }
}

代码说明

  • explode:将每行的天数序列拆分为单行数据,方便后续去重操作;
  • trim:处理输入中可能存在的空格(如示例里的"fri "),避免出现重复的"fri"和"fri ";
  • distinct:获取所有唯一的天数;
  • 广播变量:将全局唯一天数集合广播到所有Task节点,减少数据传输开销;
  • lit + cast:将广播的数组转换为Spark的Array类型列。

注意事项

如果全局唯一天数的数量极大,collect操作可能会导致Driver端内存压力过大,但此类场景下天数属于有限枚举值,该方案完全适用。

内容的提问来源于stack exchange,提问作者Didi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 21:01:04