在 Spark 中嵌套并行化?什么是正确的做法?

Jim*_*hse 2 java parallel-processing nested apache-spark

嵌套并行化?

假设我正在尝试在 Spark 中执行相当于“嵌套 for 循环”的操作。就像在常规语言中一样,假设我在内部循环中有一个例程,它以 Pi 平均 Spark 示例的方式估计 Pi(请参阅估计 Pi)

i = 1000; j = 10^6; counter = 0.0;

for ( int i =0; i < iLimit; i++)
    for ( int j=0; j < jLimit ; j++)
        counter += PiEstimator();

estimateOfAllAverages = counter / i;
Run Code Online (Sandbox Code Playgroud)

我可以在 Spark 中嵌套并行化调用吗?我正在尝试并且还没有解决问题。很乐意发布错误和代码,但我想我在问一个更具概念性的问题,关于这是否是 Spark 中的正确方法。

我已经可以并行化单个 Spark Example / Pi Estimate,现在我想这样做 1000 次以查看它是否收敛于 Pi。(这与我们试图解决的一个更大的问题有关,如果需要更接近 MVCE 的东西,我很乐意添加)

底线问题我只需要有人直接回答:这是使用嵌套并行化调用的正确方法吗?如果不是请指点一下具体的,谢谢!这是我认为正确方法的伪代码方法:

// use accumulator to keep track of each Pi Estimate result

sparkContext.parallelize(arrayOf1000, slices).map{ Function call

     sparkContext.parallelize(arrayOf10^6, slices).map{
            // do the 10^6 thing here and update accumulator with each result
    }
}

// take average of accumulator to see if all 1000 Pi estimates converge on Pi
Run Code Online (Sandbox Code Playgroud)

背景:我已经问过这个问题并得到了一个一般性的答案,但没有找到解决方案,经过一番胡思乱想后,我决定发布一个具有不同特征的新问题。我也尝试在 Spark 用户邮件列表上询问这个问题,但那里也没有骰子。在此先感谢您的帮助。

Jus*_*ony 5

这甚至是不可能的,因为它SparkContext是不可序列化的。如果你想要一个嵌套的 for 循环,那么你最好的选择是使用cartesian

val nestedForRDD = rdd1.cartesian(rdd2)
nestedForRDD.map((rdd1TypeVal, rdd2TypeVal) => {
  //Do your inner-nested evaluation code here
})
Run Code Online (Sandbox Code Playgroud)

请记住,就像双for循环一样,这是有大小成本的。