【问题标题】:How to Perform a PipelineModel.transform inside an spark.sql query如何在 spark.sql 查询中执行 PipelineModel.transform
【发布时间】:2018-09-27 01:11:35
【问题描述】:

我有一个带有列的 DataFrame“testData”:

"PRODUCT_LINE","PROFESSION","GENDER","MARITAL_STATUS"

使用测试数据执行一些预测。 我必须从其他属性预测“PRODUCT_LINE”,为此,我创建了一个 ModelRandomForest,它是一个 PipelineModel“modelrf_loaded”。

它的存储和加载,工作正常,但现在我想在 Spark SQL 中以这种方式在查询(可能使用 udf)中动态执行此操作:

val prediction3 = spark.sql("SELECT predictProductLine(*) FROM testData")

除了能够使用这样的不同语法获得我想要的结果之外:

def predictProductLine(dataframe:DataFrame):DataFrame = {
    modelrf_loaded.transform(dataframe) 
}
predictProductLine(spark.sql("SELECT * FROM testData")).show

我明确希望在单个 SQL 查询中执行操作!有没有这样使用 Spark SQL 的解决方案?

编辑:

我想出了一些可能可以像这样使用的 udf:

def udfTest(GENDER:String,AGE:Integer,MARITAL_STATUS:String,PROFESSION:String):String = {
    //case class ObjData(GENDER:String,AGE:Integer,MARITAL_STATUS:String,PROFESSION:String)
    val vardata = Seq(ObjData(GENDER,AGE,MARITAL_STATUS,PROFESSION)).toDF()
    var modelrf_loaded = PipelineModel.load("ModelRandomForest")
    val prediction = modelrf_loaded.transform(vardata)
    (prediction.first.getString(11))
}

spark.udf.register("Predict",udfTest _)
val prediction = spark.sql("SELECT Predict(GENDER,AGE,MARITAL_STATUS,PROFESSION) FROM testData")

我不是 scala/spark 方面的专家,所以也许还有另一种解决方案,但这是我脑子里唯一想到的。

在这个解决方案中存在两个问题:

  • 我需要在调用外定义类ObjData,udf内的类定义引发编译错误,所以我不能在1行中随意使用
  • 当我尝试执行 prediction.show 时,它会引发 NullPointerException,尽管它可以编译

【问题讨论】:

    标签: scala apache-spark apache-spark-sql spark-dataframe apache-spark-mllib


    【解决方案1】:

    您不能将 ML 模型与 SQL 一起使用,并且它们不能应用于 udf,因为转换方法是 (Dataset) => Dataset。您必须在此处使用标准的Dataset API。

    【讨论】:

    • 也许我可以使用 udf 来执行不同对象类型之间的更改,然后在内部进行操作,然后返回支持的类型和预测。我用一个例子更新了这个问题。感谢您的回答!
    猜你喜欢
    • 1970-01-01
    • 2022-06-27
    • 1970-01-01
    • 1970-01-01
    • 2017-07-21
    • 1970-01-01
    • 2016-02-24
    • 2020-02-21
    • 2021-06-14
    相关资源
    最近更新 更多