【问题标题】:Apache Spark's RDD[Vector] Immutability issueApache Spark 的 RDD[Vector] 不变性问题
【发布时间】:2016-04-03 12:07:47
【问题描述】:

我知道 RDD 是不可变的,因此它们的值无法更改,但我看到以下行为:

我为 FuzzyCMeans (https://github.com/salexln/FinalProject_FCM) 算法编写了一个实现,现在我正在测试它,所以我运行以下示例:

import org.apache.spark.mllib.clustering.FuzzyCMeans
import org.apache.spark.mllib.linalg.Vectors

val data = sc.textFile("/home/development/myPrjects/R/butterfly/butterfly.txt")
val parsedData = data.map(s => Vectors.dense(s.split(' ').map(_.toDouble))).cache()
> parsedData: org.apache.spark.rdd.RDD[org.apache.spark.mllib.linalg.Vector] = MapPartitionsRDD[2] at map at <console>:31

val numClusters = 2
val numIterations = 20


parsedData.foreach{ point => println(point) }
> [0.0,-8.0]
[-3.0,-2.0]
[-3.0,0.0]
[-3.0,2.0]
[-2.0,-1.0]
[-2.0,0.0]
[-2.0,1.0]
[-1.0,0.0]
[0.0,0.0]
[1.0,0.0]
[2.0,-1.0]
[2.0,0.0]
[2.0,1.0]
[3.0,-2.0]
[3.0,0.0]
[3.0,2.0]
[0.0,8.0] 

val clusters = FuzzyCMeans.train(parsedData, numClusters, numIteration
parsedData.foreach{ point => println(point) }
> 
[0.0,-0.4803333185624595]
[-0.1811743096972924,-0.12078287313152826]
[-0.06638890786148487,0.0]
[-0.04005925925925929,0.02670617283950619]
[-0.12193263222069807,-0.060966316110349035]
[-0.0512,0.0]
[NaN,NaN]
[-0.049382716049382706,0.0]
[NaN,NaN]
[0.006830134553650707,0.0]
[0.05120000000000002,-0.02560000000000001]
[0.04755220304297078,0.0]
[0.06581619798335057,0.03290809899167529]
[0.12010867103812725,-0.0800724473587515]
[0.10946638900458144,0.0]
[0.14814814814814817,0.09876543209876545]
[0.0,0.49119985188436205] 

但是我的方法怎么会改变 Immutable RDD?

顺便说一句,train 方法的签名如下:

train(数据:RDD[Vector],集群:Int,maxIterations:Int)

【问题讨论】:

  • 您能否接受其中一个答案或解释为什么这些对您不起作用,以便可以改进这些?提前致谢。

标签: scala apache-spark rdd apache-spark-mllib


【解决方案1】:

the docs 中准确描述了您所做的事情:

打印 RDD 的元素

另一个常见的习惯用法是尝试打印出 RDD 的元素 使用 rdd.foreach(println) 或 rdd.map(println)。在单机上, 这将生成预期的输出并打印所有 RDD 元素。但是,在集群模式下,输出到 stdout 被调用 执行者现在正在写入执行者的标准输出,而不是 驱动程序上的那个,所以驱动程序上的标准输出不会显示这些!到 打印驱动程序上的所有元素,可以使用 collect() 方法 首先将RDD带到驱动节点,因此: rdd.collect().foreach(println)。这可能会导致驱动程序耗尽 但是,因为 collect() 将整个 RDD 提取到一个 单机;如果您只需要打印 RDD 的几个元素,则 更安全的方法是使用 take(): rdd.take(100).foreach(println)。

因此,由于数据可以在节点之间迁移,因此无法保证 foreach 的相同输出。 RDD 是不可变的,但您应该以适当的方式提取数据,因为您的节点上没有完整的 RDD。


另一个可能的问题(不是您的情况,因为您使用的是不可变向量)是在 Point iself 中使用可变数据,这是完全不正确的,因此您将失去所有保证 - RDD 本身仍然是但是将是不可变的。

【讨论】:

    【解决方案2】:

    为了使 RDD 完全不可变,它的内容也应该是不可变的:

    scala> val m = Array.fill(2, 2)(0)
    m: Array[Array[Int]] = Array(Array(0, 0), Array(0, 0))
    
    scala> val rdd = sc.parallelize(m)
    rdd: org.apache.spark.rdd.RDD[Array[Int]] = ParallelCollectionRDD[1]
    at parallelize at <console>:23
    
    scala> rdd.collect()
    res6: Array[Array[Int]] = Array(Array(0, 0), Array(0, 0))
    
    scala> m(0)(1) = 2
    
    scala> rdd.collect()
    res8: Array[Array[Int]] = Array(Array(0, 2), Array(0, 0)) 
    

    所以因为数组是可变的,我可以改变它,因此 RDD 用新数据更新

    【讨论】:

      猜你喜欢
      • 2020-12-18
      • 2016-07-19
      • 2014-05-13
      • 2023-03-17
      • 1970-01-01
      • 2017-07-02
      • 1970-01-01
      • 1970-01-01
      • 2017-03-07
      相关资源
      最近更新 更多