【问题标题】:How add new column based on existing column in spark scala如何在 spark scala 中基于现有列添加新列
【发布时间】:2015-10-05 19:51:50
【问题描述】:

光环

我已完成在 apache spark 中使用 Mllib ALS 构建推荐,并带有输出

user | product | rating
    1 | 20 | 0.002
    1 | 30 | 0.001
    1 | 10 | 0.003
    2 | 20 | 0.002
    2 | 30 | 0.001
    2 | 10 | 0.003

但我需要根据评分排序更改数据结构,就像这样:

user | product | rating | number_rangking
    1 | 10 | 0.003 | 1
    1 | 20 | 0.002 | 2 
    1 | 30 | 0.001 | 3
    2 | 10 | 0.002 | 1
    2 | 20 | 0.001 | 2
    2 | 30 | 0.003 | 3

我该怎么做?也许任何人都可以给我一个线索...

谢谢

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    您只需要一个窗口函数,具体取决于您选择rank 或rowNumber 的详细信息

    import org.apache.spark.sql.expressions.Window
    import org.apache.spark.sql.functions.rank
    
    val w = Window.partitionBy($"user").orderBy($"rating".desc)
    
    df.select($"*", rank.over(w).alias("number_rangking")).show
    // +----+-------+------+---------------+
    // |user|product|rating|number_rangking|
    // +----+-------+------+---------------+
    // |   1|     10| 0.003|              1|
    // |   1|     20| 0.002|              2|
    // |   1|     30| 0.001|              3|
    // |   2|     10| 0.003|              1|
    // |   2|     20| 0.002|              2|
    // |   2|     30| 0.001|              3|
    // +----+-------+------+---------------+
    

    使用纯 RDD 你可以groupByKey,在本地处理和flatMap:

    rdd
      // Convert to PairRDD
      .map{case (user, product, rating) => (user, (product, rating))}
      .groupByKey 
      .flatMap{case (user, vals) => vals.toArray
        .sortBy(-_._2) // Sort by rating
        .zipWithIndex // Add index
        // Yield final values
        .map{case ((product, rating), idx) => (user, product, rating, idx + 1)}}
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-01-11
      • 1970-01-01
      • 1970-01-01
      • 2018-12-25
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多