【问题标题】:Spark Cannot call a function form .WithcolumnSpark无法调用函数表单.Withcolumn
【发布时间】:2017-05-26 07:48:02
【问题描述】:

我有一个类似下面的架构,它是 collect_list 的输出 groupby

root
|
|-- usedServiceUnits: array (nullable = true)
|    |-- element: string (containsNull = true)
|-- accumulators: array (nullable = true)
|    |-- element: string (containsNull = true)

其中的值如下所示

+----------------+
|usedServiceUnits|
+----------------+
|[180, 180, 1]   |==> this is an array of String
|[180, 180, 1]   |
+----------------+

我必须在这个字段上调用def,比如

abc.select("serviceId", "recordId", "usedServiceUnits")
.withColumn("usedServiceUnits1",lit(sumAllValuesinString($"usedServiceUnits"))

def sumAllValuesinString(inString: String): String= {
 var sum = 0
 val DELIM =','
 val a = splitString(inString,DELIM)
 for ( x <- a){
   sum += Integer.parseInt(x)
 }
 sum.toString()

}

如何调用此函数并将 sum 作为返回值并设置到我的新列 - usedServiceUnits1。对于更多功能不同的领域,我需要进行类似的计算。所以基本上我正在寻找如何将其传递给我的函数或在哪里更改?

提前感谢您的建议。

【问题讨论】:

    标签: arrays scala apache-spark


    【解决方案1】:

    这适用于字符串数组,根据您的问题使用UDF

    val getSumOf = udf((value : Seq[String]) => value.map(_.toInt).sum.toString) 
    abc.withColumn("usedServiceUnits1",udf(getSumOf($"usedServiceUnits"))
    

    希望这对你有效。

    【讨论】:

      【解决方案2】:

      据我了解您的要求,我建议您使用udf 函数

      定义udf函数为

      def sumAllValuesinString = udf((inString: mutable.WrappedArray[String]) => {
        var sum = 0
        val DELIM =','
        for ( x <- inString){
          sum += Integer.parseInt(x)
        }
        sum.toString()
      })
      

      然后使用withColumn as 调用udf 函数

      abc.withColumn("usedServiceUnits1", sumAllValuesinString($"usedServiceUnits"))
      

      我希望这是你需要的

      【讨论】:

        猜你喜欢
        • 2023-04-03
        • 2021-01-28
        • 1970-01-01
        • 2020-10-04
        • 1970-01-01
        • 1970-01-01
        • 2017-05-15
        • 2018-12-26
        • 1970-01-01
        相关资源
        最近更新 更多