【问题标题】:Spark Dataframe, add new column with function using other columnsSpark Dataframe,使用其他列添加具有功能的新列
【发布时间】:2021-12-08 00:47:26
【问题描述】:

在我的 scala 程序中,我有一个数据框 df,其中包含两列 ab(均为 Int 类型)。除此之外,我还有一个先前定义的对象obj,其中包含一些方法和属性。在这里,我想使用 obj 中的数据框和属性的当前值向我的数据框 df 添加一个新列。

例如,如果我有下面的数据框:

+---+---+
| a | b |
+---+---+
| 1 | 0 |
| 4 | 8 |
| 2 | 5 |
+---+---+

如果obj 具有属性num: Int = 10 以及方法f(a: Int, b: Int): Int = {a + b - this.num},我想使用f 创建新列c,如下所示:

+---+---+-----+
| a | b |  c  |
+---+---+-----+
| 1 | 0 | -9  |
| 4 | 8 |  2  |
| 2 | 5 | -3  |
+---+---+-----+

所以想法是:对于每一行,取列 ab 的值,并在 obj 上调用 f 方法,使用 ab 作为参数也得到然后我们存储在新列c的相应行中。我试图做这样的事情:

df = df.withColumn("c", obj.f(col("a"), col("b")))

但显然它不起作用,因为col() 返回一列而不是该列的元素。我还尝试在一个用 0 填充的新列上进行 foreach 逐行填充该列,但效果不佳。

你知道我如何在 Scala 中实现这一点吗?

谢谢。

【问题讨论】:

    标签: java scala dataframe apache-spark


    【解决方案1】:

    不使用函数也能达到同样的效果,性能会更好:

    val num = 10
    df.withColumn("c", col("a") + col("b") - lit(num))
    

    UDF 版本:

    val num = 10
    val f = (a: Int, b: Int) => {a + b - num}
    val fUDF = udf(f)
    df.withColumn("c", fUDF(col("a"), col("b")))
    

    【讨论】:

    • 这只是一个例子,我真的不得不通过这个功能。在更一般的情况下,我操作比 Int 更复杂的类型,这是行不通的。我在示例中使用 Int 使其变得简单
    • 感谢您的编辑,现在可以正常使用了!我不认为在这种情况下可以使用用户定义函数(UDF)。感谢您的提示;)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-02-14
    • 2016-08-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-25
    相关资源
    最近更新 更多