如何在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
相关产品推荐
相关产品推荐

