【发布时间】:2016-02-26 14:54:06
【问题描述】:
我有三个矩阵(A、B 和 C)作为单独的 RDD,我需要将它们划分为工作节点,作为矩阵块。我执行的操作需要更新矩阵块,但我需要在矩阵块上同步,以便两个工作节点不会同时更新同一个矩阵块。我怎样才能实现这种同步。有锁定机制吗?我对 Spark (PySpark) 很陌生。
是否可以控制 Spark 进行分区的方式,即控制将哪个块发送到哪个工作节点?
请帮忙。
【问题讨论】:
标签: apache-spark pyspark
我有三个矩阵(A、B 和 C)作为单独的 RDD,我需要将它们划分为工作节点,作为矩阵块。我执行的操作需要更新矩阵块,但我需要在矩阵块上同步,以便两个工作节点不会同时更新同一个矩阵块。我怎样才能实现这种同步。有锁定机制吗?我对 Spark (PySpark) 很陌生。
是否可以控制 Spark 进行分区的方式,即控制将哪个块发送到哪个工作节点?
请帮忙。
【问题讨论】:
标签: apache-spark pyspark
从技术上讲,这完全没有关系。 Spark 中不存在共享的、可变的状态(有人可能会争辩说accumulators 就是这种情况,但不要纠缠于此)。这意味着不存在计算可以修改共享状态并且需要任何类型的锁的情况。
这在 JVM 上稍微复杂一些,但 PySpark 架构提供了工作人员之间的完全隔离,所以除非你走出 Spark 的保险箱。如果您这样做,您有责任使用特定于上下文的方法处理冲突。
最后,如果您尝试修改数据(请不要将其与 RDD 混合),这只是一个编程错误。它可能会在 JVM 上导致一些非常丑陋的事情,但对 PySpark 应该再一次没有明显的影响(这只是实现问题而不是合同问题)。每个更改都应使用转换来表示,并且只要未另行指定(参见例如 fold 或 aggregate 系列),不应修改现有数据。
【讨论】: