【发布时间】: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