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

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()
  }
}
输入输出样例

输入样例

dateaccountidscoregenreviewed
2023-10-01acc18.5Actiontrue
2023-10-01acc17.2Comedytrue
2023-10-02acc19.0Actiontrue
2023-10-01acc26.8Dramatrue

输出样例

+---------+------+-----------------+
|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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 08:57:40