【问题标题】:Merge multiple RDD generated in loop合并循环生成的多个RDD
【发布时间】:2016-06-30 06:48:19
【问题描述】:

我在 scala 中调用一个函数,它给出一个RDD[(Long,Long,Double)] 作为它的输出。

def helperfunction(): RDD[(Long, Long, Double)]

我在代码的另一部分循环调用这个函数,我想合并所有生成的 RDD。调用函数的循环看起来像这样

for (i <- 1 to n){
    val tOp = helperfunction()
    // merge the generated tOp
}

我想做的事情类似于 StringBuilder 在你想要合并字符串时在 Java 中为你做的事情。我看过合并 RDD 的技术,主要指向使用这样的联合函数

RDD1.union(RDD2)

但这需要在合并之前生成两个 RDD。我虽然初始化了一个 var RDD1 以在 for 循环之外累积结果,但我不确定如何初始化 [(Long,Long,Double)] 类型的空白 RDD。我也是从 spark 开始的,所以我什至不确定这是否是解决这个问题的最优雅的方法。

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    您可以使用函数式编程范例来实现您想要的,而不是使用 vars:

    val rdd = (1 to n).map(x => helperFunction()).reduce(_ union _)
    

    另外,如果你仍然需要创建一个空的 RDD,你可以使用:

    val empty = sc.emptyRDD[(long, long, String)]
    

    【讨论】:

    • IIRC 在 Spark 2.0 之前,您不能将 RDD 合并到 emptyRDD。
    • 如果你必须将循环索引传递给辅助函数,你该怎么做?
    • 如果你想将循环索引传递给辅助函数,你可以这样做:val rdd = (1 to n).zipWithIndex.map{ case (x, index) =&gt; helperFunction(i) }.reduce(_ union _) 当然,在这种情况下它不是必需的,因为我们有一个整数递增集合但是你可以用任何集合替换(1 to n)
    【解决方案2】:

    您说得对,这可能不是执行此操作的最佳方式,但我们需要更多信息,了解您在每次调用帮助函数时生成新 RDD 时要完成的工作。

    您可以在循环之前定义 1 个 RDD 并为其分配一个 var,然后在您的循环中运行它。这是一个例子:

    val rdd = sc.parallelize(1 to 100)
    val rdd_tuple = rdd.map(x => (x.toLong, (x*10).toLong, x.toDouble))
    var new_rdd = rdd_tuple
    println("Initial RDD count: " + new_rdd.count())
    for (i <- 2 to 4) {
      new_rdd = new_rdd.union(rdd_tuple)
    }
    println("New count after loop: " + new_rdd.count())
    

    【讨论】:

    • 任何机构都有相同场景的JavaCode?​​span>
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-10-25
    • 1970-01-01
    • 2017-12-05
    • 2019-04-09
    • 2019-07-07
    • 2020-07-04
    相关资源
    最近更新 更多