【问题标题】:PySpark: Using column name in UDF and doing concatenation of column names based on a logicPySpark:在 UDF 中使用列名并根据逻辑连接列名
【发布时间】:2021-11-12 19:53:42
【问题描述】:

我有一个像这样的 Spark 数据框...

ID A B C D
id1 1 0 0 2
id2 0 3 0 1
id3 1 2 5 0
id4 4 0 0 1

我想要一个基于此逻辑的新数据框...

  1. 接受任何具有正值的列
  2. 连接他们的名字

结果会是这样...

ID NewColumn
id1 A,D
id2 B,D
id3 A,B,C
id4 A,D

我的努力:

A) 第一步,我想我会将整数转换为列的名称... 所以它看起来像这样......

ID A B C D
id1 A 0 0 D
id2 0 B 0 D
id3 A B C 0
id4 A 0 0 D

我正在尝试使用 UDF,但它不起作用...

def CountSelect(colname, x):
  if x>0 :
    return colname
  else:
    return ""

countUDF = UserDefinedFunction(CountSelect, T.StringType())

cols = inoutDF.columns
cols.remove("ID")

intermediateDF = inputDF.select("ID", *(countDF(c, col(c)).alias(c) for c in cols))

但它不起作用......

请大家帮忙看看?

B) 然后我将在所有列上使用字符串 concat 函数

这部分应该更简单,但如果你能将这两个逻辑组合成一个更简单的工作代码,我将非常感谢你。

非常感谢

【问题讨论】:

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


    【解决方案1】:

    我已经用 UDF 解决了这个问题,我也在这里发布。

    顺便说一句:我做了一些修改,而不是将最后一列设置为用“,”分隔的字符串 - 我创建了一个列表,这可以更好地为我的项目的后续步骤服务。

        from pyspark.sql.functions import lit, col, UserDefinedFunction, array
        import pyspark.sql.types as T
    
        def MakeOne(colname, x):
          return colname if x> 0 else None
        makeOneUDF = UserDefinedFunction(MakeOne, T.StringType())
    
        cols = inputDF.columns
        cols.remove("ID")
    
        def MakeList(arr):
          return [a for a in arr if a is not None]
        makeListUDF = UserDefinedFunction(MakeList, T.ArrayType(T.StringType()))
    
        outputDF = (inputDF.select("ID", *(makeOneUDF(lit(c), col(c)).alias(c) for c in cols)).withColumn("NewColumn", makeListUDF(array(*cols) )).select("ID", "NewColumn"))
    
    

    同样,NewColumn 是数组类型或字符串类型,它存储列名列表。

    | ID   | NewColumn|
    |------|----------|
    | id1  | [A,D]   |
    | id2  | [B,D]   |
    | id3  | [A,B,C] |
    | id4  | [A,D]   |
    

    【讨论】:

      【解决方案2】:

      这个想法是在 positive 的列中标记行并返回相应列的值。

      您可以使用reduce标记列并创建一个新的DataFrame,最后使用concat_ws形成所需的值

      @anky 提供的更简洁的解决方案

      简洁的解决方案 -

      sparkDF.withColumn("GreaterThanZero",F.concat_ws(",",*[F.when(F.col(col)>0,col) for col in to_concat]))\
      .select("id","GreaterThanZero").show()
      
      +---+---------------+
      | id|GreaterThanZero|
      +---+---------------+
      |id1|            A,D|
      |id2|            B,D|
      |id3|          A,B,C|
      |id4|            A,D|
      +---+---------------+
      

      数据准备

      input_str = """
      id1 1   0   0   2
      id2 0   3   0   1
      id3 1   2   5   0
      id4 4   0   0   1
      """.split()
      
      input_values = list(map(lambda x: x.strip() if x.strip() != 'null' else None, input_str))
      
      cols = list(map(lambda x: x.strip() if x.strip() != 'null' else None, "ID   A   B   C   D".split()))
                  
      n = len(input_values)
      n_cols = 5
      
      input_list = [tuple(input_values[i:i+n_cols]) for i in range(0,n,n_cols)]
      
      sparkDF = sql.createDataFrame(input_list, cols)
      
      sparkDF.show()
      
      +---+---+---+---+---+
      | ID|  A|  B|  C|  D|
      +---+---+---+---+---+
      |id1|  1|  0|  0|  2|
      |id2|  0|  3|  0|  1|
      |id3|  1|  2|  5|  0|
      |id4|  4|  0|  0|  1|
      +---+---+---+---+---+
      
      

      减少

      to_check = ['id','A','B','C','D']
      
      sparkDF_marked = reduce(lambda df
                      , x: df.withColumn(x,F.when(F.col(x) > 0 ,x).otherwise(None))\
                              if x != 'id' else df.withColumn(x,F.col(x)) \
                      ,to_check, sparkDF
                  )
                            
      sparkDF_marked.show()
      
      +---+----+----+----+----+
      | id|   A|   B|   C|   D|
      +---+----+----+----+----+
      |id1|   A|null|null|   D|
      |id2|null|   B|null|   D|
      |id3|   A|   B|   C|null|
      |id4|   A|null|null|   D|
      +---+----+----+----+----+
      

      连接

      to_concat = ['A','B','C','D']
      
      sparkDF_marked.select(['id',F.concat_ws(',',*to_concat).alias('GreaterThanZero')]).show()
      
      +---+---------------+
      | id|GreaterThanZero|
      +---+---------------+
      |id1|            A,D|
      |id2|            B,D|
      |id3|          A,B,C|
      |id4|            A,D|
      +---+---------------+
      

      该解决方案虽然有效,但有一些细微差别需要小心,尤其是 reduce 代码 sn-pto_checkto_concat

      to_check 可以很容易地替换为 - sparkDF.columns 用于实际数据,但请告诉我在更大数据集上的性能。

      【讨论】:

      • 是的,我在创建这个时想出了一个,我没有包括那个,因为我想打破各个步骤以便更好地理解。还添加了简洁的解决方案,谢谢您的提示
      • 这很棒。与此同时,我已经用 UDF 解决了它,并让我也发布它。顺便说一句:我进行了一些修改,而不是制作字符串 concat - 我创建了一个列表,这可以更好地为接下来的步骤服务。非常感谢您的解决方案。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-02-27
      • 2021-03-08
      • 1970-01-01
      • 2014-07-26
      • 2020-11-07
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多