【问题标题】:Spark MLLIB : Compute stddev-like value for Random Forest RegressionSpark MLLIB:计算随机森林回归的标准差值
【发布时间】:2018-10-31 22:36:32
【问题描述】:

我有一些数据想要了解“正常”行为。

使用一组有限的变量,我设法用简单的平均值来做到这一点。

df.groupBy([My_Variables]) 
  .agg(
       mean("value").alias("prediction"),
       stddev("value").alias("sigma")
  )

注意:“值”是一个双字段

我也使用随机森林算法做了同样的事情,它允许我使用更多变量。

val limit_training_set:Long = 1517439600

val trainingData = df.filter(col("datetime").cast("long")<limit_training_set)
val testData = df.filter(col("datetime").cast("long")>limit_training_set)


val assembler = new VectorAssembler()
      .setInputCols(Array(
        [My_Variables]
      ))
      .setOutputCol("features")

... // (define Indexers and Imputers)

val rf = new RandomForestRegressor()
  .setNumTrees(10) 
  .setMaxDepth(18) 
  .setLabelCol("value")
  .setFeaturesCol("features")

val pipeline = new Pipeline()
    .setStages(Array([Indexers and Imputers], assembler, rf))


val paramGrid = new ParamGridBuilder()
  .addGrid(rf.numTrees, Array(5,10))
  .addGrid(rf.maxDepth, Array(10,18)) 
  .build()

// Set up cross-validation.
val re = new RegressionEvaluator()
  .setMetricName("mae")
  .setLabelCol("value")

val tv = new TrainValidationSplit()
  .setEstimator(pipeline)
  .setEvaluator(re)
  .setEstimatorParamMaps(paramGrid)
  // 80% of the data will be used for training and the remaining 20% for validation.
  .setTrainRatio(0.8)


val model = tv.fit(trainingData)

这给了我很好的预测,但与平均方法相比,我丢失了我想要的标准差信息。

除了预测之外,还有没有办法使用随机森林计算类似 stddev 的值?或者是否有其他 ML 算法更适合这种情况?

【问题讨论】:

  • 所以你需要计算一些字段的stdev,只要你对每一行的预测......对吗?
  • 我希望有与 Mean 类似的行为。对于每组变量,我想要一个预测分数和一个“偏差”分数。

标签: scala apache-spark machine-learning apache-spark-mllib random-forest


【解决方案1】:

您可以添加一个 UnaryTransformer 来计算所需字段的分数。这将为您的行添加一个新字段:

class ScoreField(override val uid: String)
  extends UnaryTransformer[Double, Double, ScoreField]
    with DefaultParamsWritable {

  def this() = this(Identifiable.randomUID("Std"))

  final val mean: DoubleParam = new DoubleParam(this, "mean", "mean")

  final val std: DoubleParam = new DoubleParam(this, "std", "std")

  def setMean(m: Double) = set(mean, m)

  def setStd(s: Double) = set(std, s)

  override protected def createTransformFunc: Double => Double =
    v => { (v - $ { mean } / $ { std }) }

  override protected def outputDataType: DataType = DoubleType

  override def copy(extra: ParamMap): ScoreField = defaultCopy(extra)
}

object ScoreField extends DefaultParamsReadable[ScoreField] {
  def apply() : ScoreField = new ScoreField()
  override def load(path: String): ScoreField = super.load(path)
}

// Create each stage for each field 
val score = new ScoreField()
score.setInputCol("inputField")
score.setOutputCol("outputField")
// Add your mean and std four your field, this must be executed previously
score.setMean(4.0)
score.setStd(2.0)

您必须添加到您的管道中:

val pipeline = new Pipeline()
.setStages(Array([Indexers and Imputers], score, assembler, rf))

一旦您将转换调用到您的管道,您将获得带有预测和得分字段的行。

希望这会有所帮助。

【讨论】:

    猜你喜欢
    • 2016-01-28
    • 2019-04-20
    • 2013-05-06
    • 2019-12-06
    • 2016-02-21
    • 2015-12-13
    • 2017-06-13
    • 2018-06-24
    • 2016-04-24
    相关资源
    最近更新 更多