【问题标题】:Adding an index to an RDD using a shared mutable state使用共享可变状态向 RDD 添加索引
【发布时间】:2016-04-05 05:04:38
【问题描述】:

以这个简单的RDD为例说明问题:

val testRDD=sc.parallelize(List((1, 2), (3, 4), (3, 6)))

我有这个函数来帮助我实现索引:

 var sum = 0; 

 def inc(l: Int): Int = {
    sum += l
    sum 
 }

现在我想为每个元组创建 id:

val indexedRDD= testRDD.map(x=>(x._1,x._2,inc(1)));

输出RDD应该是((1,2,1), (3,4,2), (3,6,3))

但事实证明所有的值都是一样的。所有元组都取 1:

((1,2,1), (3,4,1), (3,6,1))

我哪里出错了?有没有其他方法可以达到同样的效果。

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    您正在寻找:

    def zipWithIndex(): RDD[(T, Long)]
    

    但是,请注意文档:

    请注意,某些 RDD,例如由 groupBy() 返回的 RDD,不 保证分区中元素的顺序。分配给每个的索引 因此,不保证元素,如果 RDD 为 重新评估。如果需要固定订购以保证相同 索引分配,您应该使用 sortByKey() 对 RDD 进行排序或保存 到一个文件。

    【讨论】:

    • 酷...谢谢。这个工作!但这只能回答我问题的一部分。关于为什么功能不起作用的任何想法?
    • 我有另一个 RDD 是从一个带有制表符分隔字段的文件中创建的。称之为 clickRDD 。在这种情况下,它的工作正常..val pairs = clickRDD.map(x => (x.split("\t")(0), inc(1))) ..为什么它不适用于早期的情况..在这两种情况下,我都使用了 map 函数,并且都是 RDD。
    • Spark 并行运行map。每个并行任务都有自己的sum 视图。它们可以在不同的 JVM 上运行,也可以在不同的机器上运行。根据经验,您永远不应该在 Spark 的闭包中访问可变状态。
    猜你喜欢
    • 1970-01-01
    • 2012-03-11
    • 1970-01-01
    • 1970-01-01
    • 2014-09-22
    • 1970-01-01
    • 2023-03-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多