【发布时间】:2019-03-27 11:15:40
【问题描述】:
我有一个Dataset<Row>,其中有四列,两列是非原始数据类型List<Long> and List<String>。
+------+---------------+---------------------------------------------+---------------+
| Id| value | time |aggregateType |
+------+---------------+---------------------------------------------+---------------+
|0001 | [1.5,3.4,4.5]| [1551502200000,1551502200000,1551502200000] | Sum |
+------+---------------+---------------------------------------------+---------------+
我有一个 UDF3,它接受三个参数并返回一个 Double 值 UDF3<String,List<Long>,List<String>,Double>。
所以当我调用 UDF 时,它会抛出一个异常
错误
caused by java.lang.classcastexception scala.collection.mutable.wrappedarray$ofref cannot be cast to java.lang.List
但如果我将类型更改为 String 就像 UDF3<String,String,String,Double> 一样,它不会抱怨。
抛出异常的代码
UDF3<String,List<Long>,List<String>,Double> getAggregate = new UDF3<String,List<Long>,List<String>,Double>() {
public Double call(String t1,List<Long> t2,List<String> t3) throws Exception {
//do some process to return double
return double;
}
sparkSession.udf().register("getAggregate_UDF",getAggregate, DataTypes.DoubleType);
inputDS = inputDs.withColumn("value_new",callUDF("getAggregate_UDF",col("aggregateType"),col("time"),col("value")));
将所有类型改为String后的代码
UDF3<String,String,String,Double> getAggregate = new UDF3<String,String,String,Double>() {
public Double call(String t1,String t2,String t3) throws Exception {
//code to convert t2 and t3 to List<Long> and List<String> respectively
//do some process to return double
return double;
}
sparkSession.udf().register("getAggregate_UDF",getAggregate, DataTypes.DoubleType);
inputDS = inputDs.withColumn("value_new",callUDF("getAggregate_UDF",col("aggregateType"),col("time").cast("String"),col("value").cast("String")));
上述代码有效,但需要手动转换String to List。
需要帮助
I) 如何在数据集中转换非原始数据类型List<Long> and List<String> 以克服caused by java.lang.classcastexception scala.collection.mutable.wrappedarray$ofref cannot be cast to java.lang.List
II) 如果有任何解决方法,请建议我
谢谢。
【问题讨论】:
-
您能否粘贴您的 printSchema 让我们知道实际的数据类型。
标签: apache-spark apache-spark-sql