【发布时间】:2021-01-12 13:34:51
【问题描述】:
首先我绑定到Java 1.7 和Java Spark 1.6
我有很多列和数据,但让我们按照这个简单的例子。
所以假设我有一个简单的表(DataFrame)
+----+-------+
| id| name|
+----+-------+
| 1| A|
+----+-------+
| 2| B|
+----+-------+
| 3| C|
+----+-------+
每次在每个单元格上,我都会调用自定义 udf 函数来进行所需的计算。要求之一是每次在每一行之后(或在具有某种值的每一行之后)创建并附加新的 N 行。
所以,就像:
+----+-------+
| id| name|
+----+-------+
| 1| A| --> create 1 new Row (based on the udf calculations)
+----+-------+
| 2| B| --> create 2 new Rows (based on the udf calculations)
+----+-------+
| 3| C|
+----+-------+
预期结果是:
+----+-------+
| id| name|
+----+-------+
| 1| A|
+----+-------+
| | (new)|
+----+-------+
| 2| B|
+----+-------+
| | (new)|
+----+-------+
| | (new)|
+----+-------+
| 3| C|
+----+-------+
我的误解 - 最好/正确的方法是什么?
我当前面临的问题:通过dataFrame.foreach(new Function1<Row, BoxedUnit>() {...}) Serializable error.
就我个人而言,我不确定foreach 是不是最好的方法,但我必须以某种方式迭代当前的数据帧。
此外,如果我做对了,我将始终申请 unionAll 来追加新行。
也许还有其他更好的方法可以通过Spark Sql 或将其转换为RDD 等来做到这一点。
【问题讨论】:
-
我建议从 UDF 返回一个数组,然后将数组分解为多行
-
是的,在 UDF 中,我会将计算结果保存到 temp 列中,以便在迭代期间获取它。但是对于我来说,迭代本身仍然不清楚。谢谢。
-
不需要迭代。只需做类似
df.select(col("id"),col("name"),explode(my_udf("id","name")))
标签: java apache-spark apache-spark-sql