【问题标题】:Spark: Applying UDF to Dataframe Generating new Columns based on Values in DFSpark:将 UDF 应用于 Dataframe 根据 DF 中的值生成新列
【发布时间】:2017-07-27 09:24:27
【问题描述】:

我在 Scala 中转置 DataFrame 中的值时遇到问题。我最初的DataFrame看起来像这样:

+----+----+----+----+
|col1|col2|col3|col4|
+----+----+----+----+
|   A|   X|   6|null|
|   B|   Z|null|   5|
|   C|   Y|   4|null|
+----+----+----+----+

col1col2 是类型 Stringcol3col4Int

结果应该是这样的:

+----+----+----+----+------+------+------+
|col1|col2|col3|col4|AXcol3|BZcol4|CYcol4|
+----+----+----+----+------+------+------+
|   A|   X|   6|null|     6|  null|  null|
|   B|   Z|null|   5|  null|     5|  null|
|   C|   Y|   4|   4|  null|  null|     4|
+----+----+----+----+------+------+------+

这意味着三个新列应以col1col2 和提取值的列命名。提取的值来自col2col3col5 列,具体取决于哪个值不是null

那么如何实现呢?我首先想到了这样一个UDF

def myFunc (col1:String, col2:String, col3:Long, col4:Long) : (newColumn:String, rowValue:Long) = {
    if col3 == null{
        val rowValue=col4;
        val newColumn=col1+col2+"col4";
    } else{
        val rowValue=col3;
        val newColumn=col1+col2+"col3";
     }
    return (newColumn, rowValue);
}

val udfMyFunc = udf(myFunc _ ) //needed to treat it as partially applied function

但是我怎样才能以正确的方式从数据框中调用它呢?

当然,上面的所有代码都是垃圾,可能有更好的方法。由于我只是在处理第一个代码 sn-ps 让我知道...将 Int 值与 null 进行比较已经不起作用了。

感谢任何帮助!谢谢!

【问题讨论】:

标签: scala apache-spark spark-dataframe


【解决方案1】:

有一个更简单的方法:

val df3 = df2.withColumn("newCol", concat($"col1", $"col2")) //Step 1
          .withColumn("value",when($"col3".isNotNull, $"col3").otherwise($"col4")) //Step 2
          .groupBy($"col1",$"col2",$"col3",$"col4",$"newCol") //Step 3
          .pivot("newCol") // Step 4
          .agg(max($"value")) // Step 5
          .orderBy($"newCol") // Step 6
          .drop($"newCol") // Step 7

      df3.show()

步骤如下:

  1. 添加一个新列,其中包含与 col2 连接的 col1 的内容
  2. // 添加一个新列“value”,其中包含 col3 或 col4 的非空内容
  3. GroupBy 你想要的列
  4. 以 newCol 为中心,其中包含现在将成为列标题的值
  5. 按值的最大值聚合,如果 groupBy 是每个组的单值,则该值本身就是值;或者.agg(first($"value")) 如果值恰好是字符串而不是数字类型 - max 函数只能应用于数字类型
  6. 按newCol排序,所以DF是升序的
  7. 如果您不再需要此列,请删除它,或者如果您想要一列没有空值的值,请跳过此步骤

感谢@user8371915,他首先帮助我回答了我自己的关键问题。

结果如下:

+----+----+----+----+----+----+----+
|col1|col2|col3|col4|  AX|  BZ|  CY|
+----+----+----+----+----+----+----+
|   A|   X|   6|null|   6|null|null|
|   B|   Z|null|   5|null|   5|null|
|   C|   Y|   4|   4|null|null|   4|
+----+----+----+----+----+----+----+

您可能必须尝试使用​​列标题字符串连接来获得正确的结果。

【讨论】:

    【解决方案2】:

    好的,我有一个解决方法来实现我想要的。我执行以下操作:

    (1) 我按照这个建议 Derive multiple columns from a single column in a Spark DataFrame 生成了一个包含带有 [newColumnName,rowValue] 的元组的新列

    case class toTuple(newColumnName: String, rowValue: String)
    
    def createTuple (input1:String, input2:String) : toTuple = {
        //do something fancy here
        var column:String= input1 + input2
        var value:String= input1        
        return toTuple(column, value)
    }
    
    val UdfCreateTuple = udf(createTuple _)
    

    (2) 对DataFrame应用函数

    dfNew= df.select($"*", UdfCreateTuple($"col1",$"col2").alias("tmpCol")
    

    (3) 创建具有不同值newColumnName 的数组

    val dfDistinct = dfNew.select($"tmpCol.newColumnName").distinct
    

    (4) 创建一个具有不同值的数组

    var a = dfDistinct.select($"newCol").rdd.map(r => r(0).asInstanceOf[String])
    
    var arrDistinct = a.map(a => a).collect()
    

    (5) 创建键值映射

    var seqMapping:Seq[(String,String)]=Seq()
    for (i <- arrDistinct){
        seqMapping :+= (i,i)
    }
    

    (6) 将映射应用于原始数据帧,参见。 Mapping a value into a specific column based on annother column

    val exprsDistinct = seqMapping.map { case (key, target) => 
      when($"tmpCol.newColumnName" === key, $"tmpCol.rowValue").alias(target) }
    
    val dfFinal = dfNew.select($"*" +: exprsDistinct: _*)
    

    嗯,这有点麻烦,但我可以在不知道有多少列的情况下导出一组新列,同时将值转移到那个新列中。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-06-12
      • 1970-01-01
      • 2021-04-25
      • 1970-01-01
      相关资源
      最近更新 更多