【问题标题】:Are recursive computations with Apache Spark RDD possible?可以使用 Apache Spark RDD 进行递归计算吗?
【发布时间】:2015-06-17 13:24:39
【问题描述】:

我正在使用 Scala 和 Apache Spark 开发国际象棋引擎(我需要强调我的理智不是这个问题的主题)。我的问题是 Negamax 算法本质上是递归的,当我尝试天真的方法时:

class NegaMaxSparc(@transient val sc: SparkContext) extends Serializable  {
  val movesOrdering = new Ordering[Tuple2[Move, Double]]() {
    override def compare(x: (Move, Double), y: (Move, Double)): Int =
      Ordering[Double].compare(x._2, y._2)
  }

  def negaMaxSparkHelper(game: Game, color: PieceColor, depth: Int, previousMovesPar: RDD[Move]): (Move, Double) = {
    val board = game.board

    if (depth == 0) {
      (null, NegaMax.evaluateDefault(game, color))
    } else {
      val moves = board.possibleMovesForColor(color)
      val movesPar = previousMovesPar.context.parallelize(moves)

      val moveMappingFunc = (m: Move) => { negaMaxSparkHelper(new Game(board.boardByMakingMove(m), color.oppositeColor, null), color.oppositeColor, depth - 1, movesPar) }
      val movesWithScorePar = movesPar.map(moveMappingFunc)
      val move = movesWithScorePar.min()(movesOrdering)

      (move._1, -move._2)
    }
  }

  def negaMaxSpark(game: Game, color: PieceColor, depth: Int): (Move, Double) = {
    if (depth == 0) {
      (null, NegaMax.evaluateDefault(game, color))
    } else {
      val movesPar = sc.parallelize(new Array[Move](0))

      negaMaxSparkHelper(game, color, depth, movesPar)
    }
  }
}

class NegaMaxSparkBot(val maxDepth: Int, sc: SparkContext) extends Bot {
  def nextMove(game: Game): Move = {
    val nms = new NegaMaxSparc(sc)
    nms.negaMaxSpark(game, game.colorToMove, maxDepth)._1
  }
}

我明白了:

org.apache.spark.SparkException: RDD transformations and actions can only be invoked by the driver, not inside of other transformations; for example, rdd1.map(x => rdd2.values.count() * x) is invalid because the values transformation and count action cannot be performed inside of the rdd1.map transformation. For more information, see SPARK-5063.

问题是:这个算法可以用Spark递归实现吗?如果不是,那么解决该问题的正确 Spark 方法是什么?

【问题讨论】:

    标签: scala apache-spark recursion rdd chess


    【解决方案1】:

    只有驱动程序才能在 RDD 上启动计算。原因在于,尽管 RDD “感觉”像是常规的数据集合,但在幕后它们仍然是分布式集合,因此对它们启动操作需要在所有远程从站上协调执行任务,这在大多数情况下对我们来说是隐藏的。

    因此从从站递归,即直接从从站动态启动新的分布式任务是不可能的:只有驱动器才能处理这种协调。

    这是简化您的问题的一种可能替代方法(如果我理解正确的话)。这个想法是连续构建Moves的实例,每个实例代表Move从初始状态开始的完整序列。

    Moves 的每个实例都能够将自己转换为一组Moves,每个实例对应于相同的Move 序列加上一个可能的下一个Move。

    从那里驱动程序只需连续地对Moves 进行平面映射,直到我们想要的深度,生成的 RDD[Moves] 将为我们并行执行所有操作。

    该方法的缺点是所有深度级别都保持同步,即我们必须在进入下一个级别之前计算级别 n 的所有移动(即级别 n 的 RDD[Moves])。

    下面的代码没有经过测试,它可能有缺陷,甚至没有编译,但希望它提供了一个关于如何解决问题的想法。

    /* one modification to the board */
    case class Move(from: String, to: String)
    
    case class PieceColor(color: String)
    
    /* state of the game */ 
    case class Board {
    
        // TODO
        def possibleMovesForColor(color: PieceColor): Seq[Move] = 
            Move("here", "there") :: Move("there", "over there") :: Move("there", "here") :: Nil
    
        // TODO: compute a new instance of board here, based on current + this move
        def update(move: Move): Board = new Board
    }
    
    
    /** Solution, i.e. a sequence of moves*/ 
    case class Moves(moves: Seq[Move], game: Board, color: PieceColor) {    
        lazy val score = NegaMax.evaluateDefault(game, color)
    
        /** @return all valid next Moves  */
        def nextPossibleMoves: Seq[Moves] = 
            board.possibleMovesForColor(color).map { 
                nextMove => 
                  play.copy(moves = nextMove :: play.moves, 
                            game = play.game.update(nextMove)
            } 
    
    }
    
    /** Driver code: negaMax: looks for the best next move from a give game state */
    def negaMax(sc: SparkContext, game: Board, color: PieceColor, maxDepth: Int):Moves = {
    
        val initialSolution = Moves(Seq[moves].empty, game, color)
    
        val allPlays: rdd[Moves] = 
            (1 to maxDepth).foldLeft (sc.parallelize(Seq(initialSolution))) {
            rdd => rdd.flatMap(_.nextPossibleMoves)
        }
    
        allPlays.reduce { case (m1, m2) => if (m1.score < m2.score) m1 else m2}
    
    }
    

    【讨论】:

      【解决方案2】:

      这是一个在实施方面有意义的限制,但使用起来可能会很痛苦。

      您可以尝试将递归拉到顶层,就在创建和操作 RDD 的“驱动程序”代码中?比如:

      def step(rdd: Rdd[Move], limit: Int) =
        if(0 == limit) rdd
        else {
          val newRdd = rdd.flatMap(...)
          step(newRdd, limit - 1)
        }
      

      或者,总是可以通过手动显式管理“堆栈”来将递归转换为迭代(尽管这可能会导致更繁琐的代码)。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2015-06-16
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-07-30
        • 1970-01-01
        • 2010-11-14
        相关资源
        最近更新 更多