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