【发布时间】:2018-10-27 22:01:19
【问题描述】:
我有一个如下的数据框,我正在尝试获取用户分组名称的最大值(总和)。
+-----+-----------------------------+
|name |nt_set |
+-----+-----------------------------+
|Bob |[av:27.0, bcd:29.0, abc:25.0]|
|Alice|[abc:95.0, bcd:55.0] |
|Bob |[abc:95.0, bcd:70.0] |
|Alice|[abc:125.0, bcd:90.0] |
+-----+-----------------------------+
下面是我用来为用户获取最大(总和)的 udf
val maxfunc = udf((arr: Array[String]) => {
val step1 = arr.map(x => (x.split(":", -1)(0), x.split(":", -1)(1))).groupBy(_._1).mapValues(arr => arr.map(_._2.toInt).sum).maxBy(_._2)
val result = step1._1 + ":" + step1._2
result})
当我运行 udf 时,它会抛出以下错误
val c6 = c5.withColumn("max_nt", maxfunc(col("nt_set"))).show(false)
错误:无法执行用户定义的函数($anonfun$1: (array) =>string)
我如何以更好的方式实现这一点,因为我需要在更大的数据集中做到这一点
预期的结果是
expected result:
+-----+-----------------------------+
|name |max_nt |
+-----+-----------------------------+
|Bob |abc:120.0 |
|Alice|abc:220.0 |
+-----+-----------------------------+
【问题讨论】:
标签: scala apache-spark dataframe apache-spark-sql user-defined-functions