【问题标题】:Aggregate First Grouped Item from Subsequent Items从后续项目聚合第一个分组项目
【发布时间】:2016-10-16 21:24:13
【问题描述】:

我有用户游戏会话,其中包含:用户 ID、游戏 ID、分数和玩游戏时的时间戳。

from pyspark import SparkContext
from pyspark.sql import HiveContext
from pyspark.sql import functions as F

sc = SparkContext("local")

sqlContext = HiveContext(sc)

df = sqlContext.createDataFrame([
    ("u1", "g1", 10, 0),
    ("u1", "g3", 2, 2),
    ("u1", "g3", 5, 3),
    ("u1", "g4", 5, 4),
    ("u2", "g2", 1, 1),
], ["UserID", "GameID", "Score", "Time"])

期望的输出

+------+-------------+-------------+
|UserID|MaxScoreGame1|MaxScoreGame2|
+------+-------------+-------------+
|    u1|           10|            5|
|    u2|            1|         null|
+------+-------------+-------------+

我想转换数据,以便我获得用户玩的第一场比赛的最高分以及第二场比赛的最高分(如果我还可以获得所有后续比赛的最高分,则奖励)。不幸的是,我不确定如何使用 Spark SQL。

我知道我可以按 UserID、GameID 和 agg 分组以获得最高分数和最短时间。不知道如何从那里继续。

澄清:注意MaxScoreGame1和MaxScoreGame2指的是第一个和第二个游戏用户玩家;不是 GameID。

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql pyspark-sql


    【解决方案1】:

    您可以尝试结合使用 Window 函数和 Pivot。

    1. 获取按时间排序的用户 ID 分区的每个游戏的行号。
    2. 过滤到 GameNumber 为 1 或 2。
    3. 以此为中心以获得所需的输出形状。

    不幸的是,我使用的是 scala 而不是 python,但下面的内容应该很容易转移到 python 库。

    import org.apache.spark.sql.expressions.Window
    
    // Use a window function to get row number
    val rowNumberWindow = Window.partitionBy(col("UserId")).orderBy(col("Time"))  
    
    val output = {
      df
        .select(
          col("*"),
          row_number().over(rowNumberWindow).alias("GameNumber")
        )
        .filter(col("GameNumber") <= lit(2))
        .groupBy(col("UserId"))
        .pivot("GameNumber")
        .agg(
          sum(col("Score"))
        )
    }
    
    output.show()
    
    +------+---+----+
    |UserId|  1|   2|
    +------+---+----+
    |    u1| 10|   2|
    |    u2|  1|null|
    +------+---+----+
    

    【讨论】:

    • 如果你想在输出中看到两个以上的游戏,也可以添加,只是不要过滤,枢轴会处理剩下的。
    • Window 和 row_number 成功了。我将在 PySpark 中发布我的解决方案,这有点不同。您能否验证您的代码适用于某个节目,以便我给您答案?
    • 刚刚更新了输出,还注意到我实际上在枢轴上使用了 select 而不是 groupBy,这是行不通的。对您如何获得 5 作为用户 1 的第二场比赛得分感兴趣,假设原始数据框中有错字,根据您的帖子 ("u1", "g3", 2, 2), ("u1", "g3", 5, 3),
    【解决方案2】:

    PySpark 解决方案:

    from pyspark.sql import Window
    
    rowNumberWindow = Window.partitionBy("UserID").orderBy(F.col("Time"))
    
    (df
     .groupBy("UserID", "GameID")
     .agg(F.max("Score").alias("Score"),
          F.min("Time").alias("Time"))
     .select(F.col("*"),
             F.row_number().over(rowNumberWindow).alias("GameNumber"))
     .filter(F.col("GameNumber") <= F.lit(2))
     .withColumn("GameMaxScoreCol", F.concat(F.lit("MaxScoreGame"), F.col("GameNumber")))
     .groupBy("UserID")
     .pivot("GameMaxScoreCol")  
     .agg(F.max("Score"))
    ).show()
    
    +------+-------------+-------------+
    |UserID|MaxScoreGame1|MaxScoreGame2|
    +------+-------------+-------------+
    |    u1|           10|            5|
    |    u2|            1|         null|
    +------+-------------+-------------+
    

    【讨论】:

      猜你喜欢
      • 2023-03-29
      • 1970-01-01
      • 2020-08-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-01-22
      相关资源
      最近更新 更多