【问题标题】:Training ml models on spark per partitions. Such that there will be a trained model per partition of dataframe在每个分区的 spark 上训练 ml 模型。这样每个数据框的分区都会有一个训练好的模型
【发布时间】:2020-05-13 17:35:50
【问题描述】:

如何使用 Scala 在 Spark 中对每个分区进行并行模型训练? 这里给出的解决方案在 Pyspark 中。我正在寻找scala中的解决方案。 How can you efficiently build one ML model per partition in Spark with foreachPartition?

【问题讨论】:

    标签: apache-spark apache-spark-ml


    【解决方案1】:
    1. 使用 partition col 获取不同的分区
    2. 创建一个包含 100 个线程的线程池
    3. 为每个线程创建未来对象并运行

    示例代码可能如下-

       // Get an ExecutorService 
        val threadPoolExecutorService = getExecutionContext("name", 100)
    // check https://github.com/apache/spark/blob/master/mllib/src/main/scala/org/apache/spark/ml/param/shared/HasParallelism.scala#L50
    
       val uniquePartitionValues: List[String] = ...//getDistingPartitionsUsingPartitionCol
        // Asynchronous invocation to training. The result will be collected from the futures.
        val uniquePartitionValuesFutures = uniquePartitionValues.map(partitionValue => {
          Future[Double] {
            try {
                // get dataframe where partitionCol=partitionValue
                val partitionDF = mainDF.where(s"partitionCol=$partitionValue")
              // do preprocessing and training using any algo with an input partitionDF and return accuracy
            } catch {
              ....
          }(threadPoolExecutorService)
        })
    
        // Wait for metrics to be calculated
        val foldMetrics = uniquePartitionValuesFutures.map(Await.result(_, Duration.Inf))
        println(s"output::${foldMetrics.mkString("  ###  ")}")
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-10-04
      • 2018-10-23
      • 2018-01-01
      • 2015-10-25
      • 2019-09-10
      • 2018-05-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多