【问题标题】:How to merge duplicate columns in pyspark?如何合并pyspark中的重复列?
【发布时间】:2021-09-02 19:25:52
【问题描述】:

我有一个 pyspark 数据框,其中一些列具有相同的名称。我想将所有具有相同名称的列合并到一列中。 例如,输入数据框:

如何在 pyspark 中执行此操作?任何帮助将不胜感激。

【问题讨论】:

  • dataframe 是否允许重复列?
  • 是的,由于列重命名等操作,数据框有重复列
  • 无法选择重复的列。您需要重新处理之前的处理步骤,以确保列名不重复

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


【解决方案1】:

已编辑以回答从列表合并的 OP 请求,

这是一个可重现的例子

    import pyspark.sql.functions as F

    df = spark.createDataFrame([
        ("z","a", None, None),
        ("b",None,"c", None),
        ("c","b", None, None),
        ("d",None, None, "z"),
    ], ["a","c", "c","c"])
    
    df.show()
    
    #fix duplicated column names
    old_col=df.schema.names
    running_list=[]
    new_col=[]
    i=0
    for column in old_col:
        if(column in running_list):
            new_col.append(column+"_"+str(i))
            i=i+1
        else:
            new_col.append(column)
            running_list.append(column)
    print(new_col)
    
    df1 = df.toDF(*new_col)
    
    #coalesce columns to get one column from a list

a=['c','c_0','c_1']
to_drop=['c_0','c_1']
b=[]
[b.append(df1[col]) for col in a]

#coalesce columns to get one column
df_merged=df1.withColumn('c',F.coalesce(*b)).drop(*to_drop)
   
df_merged.show()

输出:

+---+----+----+----+
|  a|   c|   c|   c|
+---+----+----+----+
|  z|   a|null|null|
|  b|null|   c|null|
|  c|   b|null|null|
|  d|null|null|   z|
+---+----+----+----+

['a', 'c', 'c_0', 'c_1']

+---+---+
|  a|  c|
+---+---+
|  z|  a|
|  b|  c|
|  c|  b|
|  d|  z|
+---+---+

【讨论】:

  • 感谢您的建议。这里的问题是我不知道没有。重复的列数以及重复的次数。那么,我们在使用coalesce的时候,应该在coalesce和drop中传入什么?
  • 实际上,我可以管理要合并到列表中的列名。例如,我有一个 python 列表 a=["'c_0","c_1","c_2"] 。现在,我怎样才能在 coalesce 中传递这个?
  • @Priyanshu 我已经编辑了代码来回答你的问题。基本上你只需要传递列表并将其解压缩并使用* 删除。如果我帮助了你,请随时将我的答案标记为已接受!
  • 非常感谢您的帮助。它现在工作正常。是的,即使我们只是通过 *a 合并,它也可以正常工作。但是两天前简单地通过 *a 内部合并是行不通的。只是想知道是否有人更新了 spark 中的合并功能。无论如何,我的问题解决了。再次感谢。
【解决方案2】:

检查下面的scala 代码。它可能会对你有所帮助。

scala> :paste
// Entering paste mode (ctrl-D to finish)

import org.apache.spark.sql.{DataFrame, SparkSession}
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._
import scala.annotation.tailrec
import scala.util.Try

implicit class DFHelpers(df: DataFrame) {
   def mergeColumns() = {
       val dupColumns = df.columns
       val newColumns = dupColumns.zipWithIndex.map(c => s"${c._1}${c._2}")
       val columns = newColumns
                        .map(c => (c(0),c))
                        .groupBy(_._1)
                        .map(c => (c._1,c._2.map(_._2)))
                        .map(c => s"""coalesce(${c._2.mkString(",")}) as ${c._1}""")
                        .toSeq
       df.toDF(newColumns:_*).selectExpr(columns:_*)
   }
}

// Exiting paste mode, now interpreting.
scala> df.show(false)
+----+----+----+----+----+----+
|a   |b   |a   |c   |a   |b   |
+----+----+----+----+----+----+
|4   |null|null|8   |null|21  |
|null|8   |7   |6   |null|null|
|96  |null|null|null|null|78  |
+----+----+----+----+----+----+
scala> df.printSchema
root
 |-- a: string (nullable = true)
 |-- b: string (nullable = true)
 |-- a: string (nullable = true)
 |-- c: string (nullable = true)
 |-- a: string (nullable = true)
 |-- b: string (nullable = true)

scala> df.mergeColumns.show(false)
+---+---+----+
|b  |a  |c   |
+---+---+----+
|21 |4  |8   |
|8  |7  |6   |
|78 |96 |null|
+---+---+----+

【讨论】:

  • 嘿,谢谢。看起来不错,但我不熟悉scala。如果有人可以为python提供类似的解决方案,那将非常有帮助
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-12-08
  • 2020-11-29
  • 2016-05-02
  • 1970-01-01
  • 1970-01-01
  • 2021-12-29
  • 1970-01-01
相关资源
最近更新 更多