【问题标题】:Convert Array to DenseVector in Spark DataFrame using Java使用 Java 将数组转换为 Spark DataFrame 中的 DenseVector
【发布时间】:2019-03-26 09:47:33
【问题描述】:

我正在运行 Spark 2.3。我想将以下DataFrame 中的features 列从ArrayType 转换为DenseVector。我在 Java 中使用 Spark。

+---+--------------------+
| id|            features|
+---+--------------------+
|  0|[4.191401, -1.793...|
| 10|[-0.5674514, -1.3...|
| 20|[0.735613, -0.026...|
| 30|[-0.030161237, 0....|
| 40|[-0.038345724, -0...|
+---+--------------------+

root
 |-- id: integer (nullable = false)
 |-- features: array (nullable = true)
 |    |-- element: float (containsNull = false)

我已经写了以下UDF,但它似乎不起作用:

private static UDF1 toVector = new UDF1<Float[], Vector>() {

    private static final long serialVersionUID = 1L;

    @Override
    public Vector call(Float[] t1) throws Exception {

        double[] DoubleArray = new double[t1.length];
        for (int i = 0 ; i < t1.length; i++)
        {
            DoubleArray[i] = (double) t1[i];
        }   
    Vector vector = (org.apache.spark.mllib.linalg.Vector) Vectors.dense(DoubleArray);
    return vector;
    }
}

我希望将以下特征提取为向量,以便对其进行聚类。

我也在注册 UDF,然后继续调用它,如下所示:

spark.udf().register("toVector", (UserDefinedAggregateFunction) toVector);
df3 = df3.withColumn("featuresnew", callUDF("toVector", df3.col("feautres")));
df3.show();  

在运行这个 sn-p 时,我遇到了以下错误:

ReadProcessData$1 不能转换为 org.apache.spark.sql.expressions。用户定义聚合函数

【问题讨论】:

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


    【解决方案1】:

    问题在于您如何在 Spark 中注册 udf。您不应使用UserDefinedAggregateFunction,它不是udf,而是用于聚合的udaf。相反,您应该做的是:

    spark.udf().register("toVector", toVector, new VectorUDT());
    

    然后要使用注册的函数,使用:

    df3.withColumn("featuresnew", callUDF("toVector",df3.col("feautres")));
    

    udf 本身应稍作调整如下:

    UDF1 toVector = new UDF1<Seq<Float>, Vector>(){
    
      public Vector call(Seq<Float> t1) throws Exception {
    
        List<Float> L = scala.collection.JavaConversions.seqAsJavaList(t1);
        double[] DoubleArray = new double[t1.length()]; 
        for (int i = 0 ; i < L.size(); i++) { 
          DoubleArray[i]=L.get(i); 
        } 
        return Vectors.dense(DoubleArray); 
      } 
    };
    

    请注意,在 Spark 2.3+ 中,您可以创建一个可以直接调用的 scala 样式 udf。来自这个answer:

    UserDefinedFunction toVector = udf(
      (Seq<Float> array) -> /* udf code or method to call */, new VectorUDT()
    );
    
    df3.withColumn("featuresnew", toVector.apply(col("feautres")));
    

    【讨论】:

    • @BdEngineer:对于 Spark 中的机器学习,向量(DenseVector、SparseVector)用于输入而不是数组。也可能有其他用例。
    猜你喜欢
    • 2017-08-10
    • 2017-03-25
    • 2016-12-25
    • 1970-01-01
    • 2017-05-09
    • 1970-01-01
    • 2022-01-08
    • 1970-01-01
    • 2017-03-17
    相关资源
    最近更新 更多