【问题标题】:Transposing a Spark DataFrame from row to column in PySpark and appending it with another DataFrame在 PySpark 中将 Spark DataFrame 从行转换为列,并将其附加到另一个 DataFrame
【发布时间】:2020-02-20 01:50:09
【问题描述】:

我在 PySpark avg_length_df 中有一个 Spark DataFrame,看起来像 -

+----------------+---------+----------+-----------+---------+-------------+----------+
|       id       |        x|         a|          b|        c|      country|     param|
+----------------+---------+----------+-----------+---------+-------------+----------+
|            40.0|      9.0|     5.284|      5.047|    6.405|         13.0|avg_length|
+----------------+---------+----------+-----------+---------+-------------+----------+

我想把它从行转置到列,这样它就变成了 -

+----------+
|avg_length|
+----------+
|      40.0|
|       9.0|
|     5.284|
|     5.047|
|     6.405|
|      13.0|
+----------+

接下来,我有第二个 DataFrame df2:

+----------------+------+
|       col_names|dtypes|
+----------------+------+
|              id|string|
|               x|   int|
|               a|string|
|               b|string|
|               c|string|
|         country|string|
+----------------+------+

我想在df2 中创建一个列avg_length,等于上面的转置DataFrame。所以预期的输出看起来像:

+----------------+------+----------+
|       col_names|dtypes|avg_length|
+----------------+------+----------+
|              id|string|      40.0|
|               x|   int|       9.0|
|               a|string|     5.284|
|               b|string|     5.047|
|               c|string|     6.405|
|         country|string|      13.0|
+----------------+------+----------+

如何完成这2个操作?

【问题讨论】:

    标签: python dataframe apache-spark pyspark transpose


    【解决方案1】:
    >>> from pyspark.sql import *
    #Input DataFrame
    >>> df.show()
    +----+---+-----+-----+-----+-------+----------+
    |  id|  x|    a|    b|    c|country|     param|
    +----+---+-----+-----+-----+-------+----------+
    |40.0|9.0|5.284|5.047|6.405|   13.0|avg_length|
    +----+---+-----+-----+-----+-------+----------+
    
    >>> avgDF  = df.groupBy(df["id"],df["x"],df["a"],df["b"],df["c"],df["country"]).pivot("param").agg(concat_ws("",collect_list(to_json(struct("id","x","a","b","c","country"))))).drop("id","x","a","b","c","country")
    >>> avgDF.show(2,False)
    +----------------------------------------------------------------------------+
    |avg_length                                                                  |
    +----------------------------------------------------------------------------+
    |{"id":"40.0","x":"9.0","a":"5.284","b":"5.047","c":"6.405","country":"13.0"}|
    +----------------------------------------------------------------------------+
    
    >>> finalDF = avgDF.withColumn("value", explode(split(regexp_replace(col("avg_length"),"""[\\{ " \\}]""",""),","))).withColumn("avg_length", split(col("value"), ":")[1]).withColumn("col_names", split(col("value"), ":")[0]).drop("value")
    >>> finalDF.show(10,False)
    +----------+---------+
    |avg_length|col_names|
    +----------+---------+
    |40.0      |id       |
    |9.0       |x        |
    |5.284     |a        |
    |5.047     |b        |
    |6.405     |c        |
    |13.0      |country  |
    +----------+---------+
    
    #other dataframe
    >>> df2.show()
    +---------+------+
    |col_names|dtypes|
    +---------+------+
    |       id|string|
    |        x|   int|
    |        a|string|
    |        b|string|
    |        c|string|
    |  country|string|
    +---------+------+
    
    >>> df2.join(finalDF,"col_names").show(10,False)
    +---------+------+----------+
    |col_names|dtypes|avg_length|
    +---------+------+----------+
    |id       |string|40.0      |
    |x        |int   |9.0       |
    |a        |string|5.284     |
    |b        |string|5.047     |
    |c        |string|6.405     |
    |country  |string|13.0      |
    +---------+------+----------+
    

    【讨论】:

      【解决方案2】:

      以下是在 pyspark 中转置数据帧 (RDD) 的代码。

      import numpy as np
      from pyspark.sql import SQLContext
      from pyspark.sql.functions import lit
      
      dt1 = {'avg_length':[40.0, 9.0, 5.284, 5.047, 6.405, 13.0]}
      dt = sc.parallelize([ (k,) + tuple(v[0:]) for k,v in dt1.items()]).toDF()
      dt.show()
      
      
      #--- Transpose Code ---
      
      
      # Grad data from first columns, since it will be transposed to new column headers
      new_header = [i[0] for i in dt.select("_1").rdd.map(tuple).collect()]
      
      # Remove first column from dataframe
      dt2 = dt.select([c for c in dt.columns if c not in ['_1']])
      
      # Convert DataFrame to RDD
      rdd = dt2.rdd.map(tuple)
      
      # Transpose Data
      rddT1 = rdd.zipWithIndex().flatMap(lambda (x,i): [(i,j,e) for (j,e) in enumerate(x)])
      rddT2 = rddT1.map(lambda (i,j,e): (j, (i,e))).groupByKey().sortByKey()
      rddT3 = rddT2.map(lambda (i, x): sorted(list(x), cmp=lambda (i1,e1),(i2,e2) : cmp(i1, i2)))
      rddT4 = rddT3.map(lambda x: map(lambda (i, y): y , x))
      
      # Convert back to DataFrame (along with header)
      df = rddT4.toDF(new_header)
      
      df.show()
      

      转置后,您可以简单地合并两个数据帧。 我希望这会有所帮助。

      【讨论】:

      • 为什么对这个答案投反对票,代码不起作用?
      • 请不要在没有任何适当评论和逻辑的情况下投票。
      • 因为您从其他网站复制代码并粘贴到这里,没有任何解释。此代码不起作用。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-09-25
      • 1970-01-01
      • 2017-03-17
      • 2017-02-20
      • 2022-09-24
      相关资源
      最近更新 更多