【问题标题】:spark flatten records using a key column火花使用键列展平记录
【发布时间】:2017-10-19 20:10:49
【问题描述】:

我正在尝试使用 spark/Scala API 实现逻辑以展平记录。我正在尝试使用地图功能。

你能帮我用最简单的方法解决这个问题吗?

假设,对于给定的键,我需要 3 个进程代码

输入数据框-->

Keycol|processcode
John  |1
Mary  |8
John  |2
John  |4
Mary  |1
Mary  |7

================================

输出数据框-->

Keycol|processcode1|processcode2|processcode3
john  |1           |2           |4
Mary  |8           |1           |7

【问题讨论】:

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


    【解决方案1】:

    假设每个 Keycol 的行数相同,一种方法是将 processcode 聚合到每个 Keycol 的数组中,然后扩展到各个列:

    val df = Seq(
      ("John", 1),
      ("Mary", 8),
      ("John", 2),
      ("John", 4),
      ("Mary", 1),
      ("Mary", 7)
    ).toDF("Keycol", "processcode")
    
    val df2 = df.groupBy("Keycol").agg(collect_list("processcode").as("processcode"))
    
    val numCols = df2.select( size(col("processcode")) ).as[Int].first
    val cols = (0 to numCols - 1).map( i => col("processcode")(i) )
    
    df2.select(col("Keycol") +: cols: _*).show
    
    +------+--------------+--------------+--------------+
    |Keycol|processcode[0]|processcode[1]|processcode[2]|
    +------+--------------+--------------+--------------+
    |  Mary|             8|             1|             7|
    |  John|             1|             2|             4|
    +------+--------------+--------------+--------------+
    

    【讨论】:

    • 感谢 Leo 下面一行解决了我的问题val df2 = df.groupBy("Keycol").agg(collect_list("processcode").as("processcode")) 感谢您的快速帮助。跨度>
    【解决方案2】:

    几种替代方法。

    SQL

    df.createOrReplaceTempView("tbl")
    
    val q = """
    select keycol,
           c[0] processcode1,
           c[1] processcode2,
           c[2] processcode3
      from (select keycol, collect_list(processcode) c
              from tbl
            group by keycol) t0
    """
    
    sql(q).show
    

    结果

    scala> sql(q).show
    +------+------------+------------+------------+
    |keycol|processcode1|processcode2|processcode3|
    +------+------------+------------+------------+
    |  Mary|           1|           7|           8|
    |  John|           4|           1|           2|
    +------+------------+------------+------------+
    

    PairRDDFunctions (groupByKey) + mapPartitions

    import org.apache.spark.sql.Row
    val my_rdd = df.map{ case Row(a1: String, a2: Int) => (a1, a2)
                       }.rdd.groupByKey().map(t => (t._1, t._2.toList))
    
    def f(iter: Iterator[(String, List[Int])]) : Iterator[Row] = {
      var res = List[Row]();
      while (iter.hasNext) {
        val (keycol: String, c: List[Int]) = iter.next    
        res = res ::: List(Row(keycol, c(0), c(1), c(2)))
      }
      res.iterator
    }
    
    import org.apache.spark.sql.types.{StringType, IntegerType, StructField, StructType}
    val schema = new StructType().add(
                 StructField("Keycol", StringType, true)).add(
                 StructField("processcode1", IntegerType, true)).add(
                 StructField("processcode2", IntegerType, true)).add(
                 StructField("processcode3", IntegerType, true))
    
    spark.createDataFrame(my_rdd.mapPartitions(f, true), schema).show
    

    结果

    scala> spark.createDataFrame(my_rdd.mapPartitions(f, true), schema).show
    +------+------------+------------+------------+
    |Keycol|processcode1|processcode2|processcode3|
    +------+------------+------------+------------+
    |  Mary|           1|           7|           8|
    |  John|           4|           1|           2|
    +------+------------+------------+------------+
    

    请记住,在所有情况下,除非明确指定,否则流程代码列中的值顺序未确定

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-04-10
      • 2021-07-23
      • 2016-04-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多