Spark Scala:为账号按每日观看流派生成权重列的实现问询
嘿,我给你整理了一个基于Spark Scala的可行方案,完美解决你要生成每日流派权重序列的需求,还避开了transpose的问题——毕竟Scala Spark里确实没这个函数,咱们换个思路用Spark原生API搞定:
实现思路拆解
核心思路是先补全所有可能的(accountid, date, genre)组合,确保每个账号的每个日期都覆盖所有流派,再标记是否有观看记录,最后按账号和流派分组生成有序的权重序列:
- 第一步:提取所有唯一日期(按顺序)和流派,生成基准笛卡尔积,保证没有遗漏的日期或流派组合
- 第二步:将原始数据与基准组合左连接,标记每个组合是否存在观看记录
- 第三步:用窗口函数按日期排序后聚合,生成随日期递增的权重列表(观看则0.2,否则0.0)
完整示例代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window object GenreWeightGenerator { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("GenreWeightGenerator") .master("local[*]") // 生产环境请移除该行 .getOrCreate() // 模拟输入数据 val inputDF = spark.createDataFrame(Seq( ("2023-10-01", "acc1", 8.5, "Action", true), ("2023-10-01", "acc1", 7.2, "Comedy", true), ("2023-10-02", "acc1", 9.0, "Action", true), ("2023-10-01", "acc2", 6.8, "Drama", true) )).toDF("date", "accountid", "score", "genre", "viewed") // 1. 获取所有唯一日期(按顺序)和流派 val allSortedDates = inputDF.select("date").distinct().orderBy("date") val allGenres = inputDF.select("genre").distinct() // 2. 生成所有(accountid, date, genre)的基准组合 val accountDateGenreBase = inputDF.select("accountid").distinct() .crossJoin(allSortedDates) .crossJoin(allGenres) // 3. 左连接原始数据,标记是否有观看记录 val joinedDF = accountDateGenreBase.join( inputDF.select("accountid", "date", "genre", "viewed"), Seq("accountid", "date", "genre"), "left_outer" ).withColumn("has_viewed", when(col("viewed").isNotNull, true).otherwise(false)) // 4. 按账号+流派分组,用窗口函数按日期排序生成权重序列 val windowSpec = Window.partitionBy("accountid", "genre").orderBy("date") val sortedWeightDF = joinedDF .withColumn("weight", when(col("has_viewed"), 0.2).otherwise(0.0)) .withColumn("weight_sequence", collect_list("weight").over(windowSpec)) // 5. 提取每个分组的完整权重序列(取分组内最后一条记录的序列) val resultDF = sortedWeightDF.groupBy("accountid", "genre") .agg(last("weight_sequence").alias("weight_sequence")) // 打印结果 resultDF.show(false) spark.stop() } }
输入输出样例
输入样例
| date | accountid | score | genre | viewed |
|---|---|---|---|---|
| 2023-10-01 | acc1 | 8.5 | Action | true |
| 2023-10-01 | acc1 | 7.2 | Comedy | true |
| 2023-10-02 | acc1 | 9.0 | Action | true |
| 2023-10-01 | acc2 | 6.8 | Drama | true |
输出样例
+---------+------+-----------------+ |accountid|genre |weight_sequence | +---------+------+-----------------+ |acc1 |Action|[0.2, 0.2] | |acc1 |Comedy|[0.2, 0.0] | |acc1 |Drama |[0.0, 0.0] | |acc2 |Action|[0.0] | |acc2 |Comedy|[0.0] | |acc2 |Drama |[0.2] | +---------+------+-----------------+
方案优势说明
- 完全基于Spark分布式API,适合大数据场景,避免了单机字典存储的内存瓶颈
- 通过基准笛卡尔积补全所有组合,确保没有遗漏的日期或流派
- 用窗口函数的
collect_list保证权重序列严格按日期递增排序,符合需求 - 不需要依赖第三方库或自定义UDF,原生API即可实现
内容的提问来源于stack exchange,提问作者Masterbuilder
相关产品推荐
相关产品推荐

