【问题标题】:Not able to call a UDF with callUDF() - Spark Java无法使用 callUDF() 调用 UDF - Spark Java
【发布时间】:2017-08-02 13:58:54
【问题描述】:

我在注册后尝试使用 callUDF 调用 udf。但是,函数 validateNumber() 没有被调用。

代码如下所示:

public Dataset<Row> sampleCallUdf(Dataset<Row> dataset) {

    UDF2<Long, Long, String> validateNumber = (UDF2<Long, Long, String>) SampleClass::validateNumber;
    UDFRegistration udfRegister = CONFIG.getSparkSession().udf();
    udfRegister.register("validateNumber", validateNumber, DataTypes.StringType);

    return dataset.withColumn("rejection_reason",
                    coalesce(
                            callUDF("validateNumber", column("cookie"), column("session"))));
    }

    public static String validateNumber(Long cookie, Long session) {
           System.out.println("Into validateNumber function");
           if(cookie != 0){
             return "correct";
           }else{
             return "incorrect";
           }
    }

我正在尝试的输入是:

 Dataset<Row> input = spark().createDataFrame(Arrays.asList(
                RowFactory.create("28/05/2017 00:12:34", 0L, -2864001245604480000L, "abc" ,"90.202.190.106", 123, "abc", "xyz", "mno"),
                RowFactory.create("28/05/2017 00:12:34", 2345678L, 2864001245604480000L, "abc" ,"90.202.190.106", 123, "abc", "xyz", "mno")), TEMP_TABLE);

问题是,它甚至没有在 validateNumber() 函数中打印 sysout 语句。

【问题讨论】:

  • 对我来说工作得很好。你能检查一下你的数据集中的值吗?
  • @abaghel - 它进入了 validateNumber() 吗?
  • 或者,如果您可以让我知道您使用的是什么输入。

标签: java apache-spark apache-spark-sql user-defined-functions


【解决方案1】:

请在下面找到示例程序。

public class SparkUDF {
    public static void main(String[] args) throws Exception {
    SparkSession spark = SparkSession
            .builder()
            .appName("SparkUDF")
            .master("local[*]") 
            .getOrCreate();
    //data
    List<Tuple2<Long, Long>> inputList = new ArrayList<Tuple2<Long, Long>>();
    inputList.add(new Tuple2<Long, Long>(111l, 10011l));
    inputList.add(new Tuple2<Long, Long>(0l, 20022l));
    //Dataset
    Dataset<Row> ds = spark.createDataset(inputList, Encoders.tuple(Encoders.LONG(), Encoders.LONG())).toDF("cookie", "session");
    //udf
    UDF2<Long, Long, String> validateNumber = (UDF2<Long, Long, String>) SparkUDF::validateNumber;
    spark.udf().register("validateNumber", validateNumber, DataTypes.StringType);
    Dataset<Row> ds1 = ds.withColumn("rejection_reason",coalesce(callUDF("validateNumber", col("cookie"), col("session"))));
    ds1.show();
    spark.stop();
}

public static String validateNumber(Long cookie, Long session) {
    if (cookie != 0) {
        return "correct";
    } else {
        return "incorrect";
    }
  }
}

你会得到如下输出。

+------+-------+----------------+
|cookie|session|rejection_reason|
+------+-------+----------------+
|   111|  10011|         correct|
|     0|  20022|       incorrect|
+------+-------+----------------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-02-23
    • 2016-05-22
    • 2023-01-15
    • 1970-01-01
    • 1970-01-01
    • 2020-10-20
    • 2019-03-13
    • 1970-01-01
    相关资源
    最近更新 更多