【问题标题】:spark data frame converting row values into column name火花数据框将行值转换为列名
【发布时间】:2020-01-03 18:02:36
【问题描述】:

使用 spark 数据框我需要将行值转换为列并按用户 ID 分区并创建一个 csv 文件。


val someDF = Seq(
  ("user1", "math","algebra-1","90"),
  ("user1", "physics","gravity","70"),
  ("user3", "biology","health","50"),
  ("user2", "biology","health","100"),
  ("user1", "math","algebra-1","40"),
  ("user2", "physics","gravity-2","20")
).toDF("user_id", "course_id","lesson_name","score")

someDF.show(false)

+-------+---------+-----------+-----+
|user_id|course_id|lesson_name|score|
+-------+---------+-----------+-----+
|  user1|     math|  algebra-1|   90|
|  user1|  physics|    gravity|   70|
|  user3|  biology|     health|   50|
|  user2|  biology|     health|  100|
|  user1|     math|  algebra-1|   40|
|  user2|  physics|  gravity-2|   20|
+-------+---------+-----------+-----+

val result = someDF.groupBy("user_id", "course_id").pivot("lesson_name").agg(first("score"))

result.show(false)

+-------+---------+---------+-------+---------+------+
|user_id|course_id|algebra-1|gravity|gravity-2|health|
+-------+---------+---------+-------+---------+------+
|  user3|  biology|     null|   null|     null|    50|
|  user1|     math|       90|   null|     null|  null|
|  user2|  biology|     null|   null|     null|   100|
|  user2|  physics|     null|   null|       20|  null|
|  user1|  physics|     null|     70|     null|  null|
+-------+---------+---------+-------+---------+------+


使用上面的代码,我可以将行值(课程名称)转换为列名。 但我需要将输出保存在 csv 中的 course_wise

在 csv 中的预期输出应如下所示。

biology.csv // Expected Output

+-------+---------+------+
|user_id|course_id|health|
+-------+---------+------+
|  user3|  biology|  50  |
|  user2|  biology| 100  |
+-------+---------+-------

physics.csv // Expected Output

+-------+---------+---------+-------
|user_id|course_id|gravity-2|gravity|
+-------+---------+---------+-------+
|  user2|  physics|  50     |  null |
|  user1|  physics| 100     |  70   | 
+-------+---------+---------+-------+

**注意:csv 中的每门课程都应仅包含其特定的课程名称,并且不应包含任何不相关的课程课程名称。

实际上在 csv 中我可以在下面的格式中生成**

result.write
  .partitionBy("course_id")
  .mode("overwrite")
  .format("com.databricks.spark.csv")
  .option("header", "true")
  .save(somepath)


例如:

biology.csv // Wrong output, Due to it is containing non-relevant course lesson's(algebra-1,gravity-2,algebra-1)
+-------+---------+---------+-------+---------+------+
|user_id|course_id|algebra-1|gravity|gravity-2|health|
+-------+---------+---------+-------+---------+------+
|  user3|  biology|     null|   null|     null|    50|
|  user2|  biology|     null|   null|     null|   100|
+-------+---------+---------+-------+---------+------+

谁能帮忙解决这个问题?

【问题讨论】:

    标签: scala dataframe apache-spark apache-spark-sql spark-streaming


    【解决方案1】:

    在转之前按课程过滤:

    val result = someDF.filter($"course_id" === "physics").groupBy("user_id", "course_id").pivot("lesson_name").agg(first("score"))
    
    +-------+---------+-------+---------+
    |user_id|course_id|gravity|gravity-2|
    +-------+---------+-------+---------+
    |user2  |physics  |null   |20       |
    |user1  |physics  |70     |null     |
    

    +-------+---------+-------+---------+

    【讨论】:

    • 为什么要硬编码“物理”?如果我有大约 1000 个 course_id 怎么办?在这种情况下我们如何处理?
    • 我想创建一个课程明智的 csv 文件?假设在这种情况下我有大约 1000 门课程,而在这种情况下,每门课程都有大约 5 节课,你如何处理? @安德鲁
    【解决方案2】:

    我假设您的意思是您希望通过 course_id 将数据保存到单独的目录中。你可以使用这种方法。

    scala> val someDF = Seq(
    ("user1", "math","algebra-1","90"),
    ("user1", "physics","gravity","70"),
    ("user3", "biology","health","50"),
    ("user2", "biology","health","100"),
    ("user1", "math","algebra-1","40"),
    ("user2", "physics","gravity-2","20")
    ).toDF("user_id", "course_id","lesson_name","score")
    
    
    scala> val result = someDF.groupBy("user_id", "course_id").pivot("lesson_name").agg(first("score"))
    
    scala>     val eventNames = result.select($"course_id").distinct().collect() 
    var eventlist =eventNames.map(x => x(0).toString)
    
    
    
    for (eventName <- eventlist) {
    val course = result.where($"course_id" === lit(eventName))
    //remove null column
    
    val row = course
    .select(course.columns.map(c => when(col(c).isNull, 0).otherwise(1).as(c)): _*)
    .groupBy().max(course.columns.map(c => c): _*)
    .first
    
    val colKeep = row.getValuesMap[Int](row.schema.fieldNames)
    .map{c => if (c._2 == 1) Some(c._1) else None }
    .flatten.toArray
    
    
    var final_df = course.select(row.schema.fieldNames.intersect(colKeep)
    .map(c => col(c.drop(4).dropRight(1))): _*)
    
    
    final_df.show()
    
    final_df.coalesce(1).write.mode("overwrite").format("csv").save(s"${eventName}")
    }
    
    
    +-------+---------+------+
    |user_id|course_id|health|
    +-------+---------+------+
    |  user3|  biology|    50|
    |  user2|  biology|   100|
    +-------+---------+------+
    
    +-------+---------+-------+---------+
    |user_id|course_id|gravity|gravity-2|
    +-------+---------+-------+---------+
    |  user2|  physics|   null|       20|
    |  user1|  physics|     70|     null|
    +-------+---------+-------+---------+
    
    +-------+---------+---------+
    |user_id|course_id|algebra-1|
    +-------+---------+---------+
    |  user1|     math|       90|
    +-------+---------+---------+
    

    如果它解决了您的目的,请接受 answer.HAppy Hadoop

    【讨论】:

    • 如果我有另一个列(batchid)和 courseid 怎么办?假设我想同时保存 courseid 和 batchid ? val someDF = Seq( ("user1", "math","algebra-1","90","b1"), ("user1", "physics","gravity","70","b1"), ("user3", "biology","health","50","b2"), ("user2", "biology","health","100","b2"), ("user1", "math","algebra-1","40","b1"), ("user2", "physics","gravity-2","20","b3") ).toDF("user_id", "course_id","lesson_name","score","batch_id")
    • 我想在 batchid 和 courseid 上创建一个 csv 文件,输出应该是这样的+-------+---------+---------+-------+---------+------+---------+ |user_id|course_id|algebra-1|gravity|gravity-2|health|course_id| +-------+---------+---------+-------+---------+------+---------+ | user2| physics| null| null| 20| null|b3. | ----------------------------------------------------------------
    • 这不会影响性能吗?当我们有巨大的 courseid 和 batchid 时?
    • 这是一个完全不同的问题,所以为此添加一个新问题。如果它解决了您的问题,请接受答案。
    • 并添加一个具有预期输出的适当数据框。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-05-15
    • 1970-01-01
    • 1970-01-01
    • 2022-08-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多