【问题标题】:Apache Spark RDD - not updatingApache Spark RDD - 不更新
【发布时间】:2023-03-17 21:07:01
【问题描述】:

我创建了一个包含 Vector 的 PairRDD。

var newRDD = oldRDD.mapValues(listOfItemsAndRatings => Vector(Array.fill(2){math.random}))

稍后我会更新 RDD:

newRDD.lookup(ratingObject.user)(0) += 0.2 * (errorRate(rating) * myVector)

然而,虽然它输出了一个更新的 Vector(如控制台所示),但当我下一次调用 newRDD 时,我可以看到 Vector 值已更改。通过测试,我得出的结论是它已经变成了math.random 给出的东西——因为每次我调用newRDD 时,Vector 都会发生变化。我知道有一个谱系图,也许这与它有关。我需要将 RDD 中保存的 Vector 更新为新值,并且我需要重复执行此操作。

谢谢。

【问题讨论】:

  • newRDD 是一个 RDD,根据定义是不可变的。我认为您无法更改 RDD 中 Vector 内的值。

标签: scala apache-spark rdd


【解决方案1】:

RDD 是不可变结构,用于在集群上分布对数据的操作。 您在此处观察到的行为有两个因素在起作用:

可以每次都计算 RDD 沿袭。在这种情况下,这意味着对 newRDD 的操作可能会触发沿袭计算,因此应用Vector(Array.fill(2){math.random}) 转换并每次都会产生新值。可以使用cache 打破沿袭,在这种情况下,转换的值将在第一次应用后保存在内存和/或磁盘中。 这导致:

val randomVectorRDD = oldRDD.mapValues(listOfItemsAndRatings => Vector(Array.fill(2){math.random}))
randomVectorRDD.cache()

需要进一步考虑的第二个方面是现场突变:

newRDD.lookup(ratingObject.user)(0) += 0.2 * (errorRate(rating) * myVector)

虽然这可能在单台机器上工作,因为所有 Vector 引用都是本地的,但它不会扩展到集群,因为查找引用将被序列化并且不会保留突变。因此,它承担了为什么要使用 Spark 的问题。

要在 Spark 上实现,该算法需要重新设计,以便用转换而不是准时查找/突变来表达。

【讨论】:

    猜你喜欢
    • 2014-05-13
    • 1970-01-01
    • 2017-03-07
    • 1970-01-01
    • 1970-01-01
    • 2015-06-15
    • 2016-12-06
    • 2016-04-21
    • 2016-04-03
    相关资源
    最近更新 更多