【问题标题】:How to pass Row in UDF?如何在 UDF 中传递 Row?
【发布时间】:2018-12-16 19:20:24
【问题描述】:

我正在用 Java 编写 UDF。

我想对DateSet<Row> 执行更复杂的操作。为了那个原因 我想我需要将 DataSet<Row> 作为 UDF 的输入传递并返回输出。这是我的代码:

 UDF1<Dataset<Row>,String> myUDF = new UDF1<Dataset<Row>,String>() {
            public String call(Dataset<Row> input) throws Exception {
                System.out.println(input);
                return "test";
            }
            };

           // Register the UDF with our SQLContext
            spark.udf().register("myUDF", myUDF, DataTypes.StringType); {

但是当我尝试使用 myUDF 时,似乎 callUDF 函数只接受 Column 而不是 DataSet&lt;Row&gt;

谁能帮助我将DataSet&lt;Row&gt; 作为输入参数传递给UDF?有没有其他方法可以在 Spark SQL 中调用我的 UDF?

【问题讨论】:

  • 我已经检查过了 那并不能解决我的问题。这是在 Scala 中实现的。我正在寻找 java 中的东西。
  • 真的没有太大区别。你需要struct(all columns go here)
  • @user10465355 这可能是一个解决方案,但会改变问题的语义(即转换数据集)。

标签: java apache-spark apache-spark-sql


【解决方案1】:

但是当我尝试使用 myUDF 时,似乎 callUDF 函数只接受列而不接受数据集,任何人都可以帮助我如何将数据集作为 UDF 中的输入参数传递。有没有其他方法可以在 Spark SQL 中调用我的 UDF

这里有几个问题。

首先,UDF 是一个与(内部值)Columns 一起使用的函数。从某种意义上说,您可以使用struct 函数来组合所需的列,以假装您使用整个数据集。

但是,如果您想使用整个数据集,您确实需要一个简单地接受数据集的纯 Java/Scala 方法。 Spark 对此无能为力。它只是一个 Java/Scala 编程。

但是有一个非常好的方法,我看不到太多用处,即Dataset.transform

transform[U](t: (Dataset[T]) ⇒ Dataset[U]): Dataset[U] 用于链接自定义转换的简洁语法。

这允许链接接受 Dataset 的方法,这使得代码非常易读(并且看起来正是您想要的)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-07-22
    • 2017-08-13
    • 2017-07-21
    • 2015-09-13
    • 1970-01-01
    • 2017-12-19
    • 2020-03-31
    相关资源
    最近更新 更多