【问题标题】:Append entry to RDD using a loop [duplicate]使用循环将条目附加到 RDD [重复]
【发布时间】:2018-01-29 11:59:33
【问题描述】:

我试图在循环的每次迭代中将一个条目附加到现有的 RDD。到目前为止,我的代码是:

var newY = sc.emptyRDD[MatrixEntry]
for (j <- 0 until 8000) {
  var arrTmp = Array(MatrixEntry(j, j, 1))
  var rddTmp = sc.parallelize(arrTmp)
  newY = newY.union(rddTmp)
}

进行这 8000 次迭代时,当我尝试从该 RDD 中取 (10) 时出现错误,但如果我尝试使用较小的数字,一切都可以。 错误Exception in thread "main" java.lang.StackOverflowError at scala.collection.TraversableLike$class.builder$1(TraversableLike.scala:229) at scala.collection.TraversableLike$class.map(TraversableLike.scala:233) at scala.collection.immutable.List.map(List.scala:296) at org.apache.spark.rdd.UnionRDD.getPartitions(UnionRDD.scala:84) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:252) at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:250) at scala.Option.getOrElse(Option.scala:121)

帮助?

【问题讨论】:

    标签: apache-spark rdd


    【解决方案1】:

    您遇到的问题与Stackoverflow due to long RDD Lineage 重复,但您的代码根本不应该存在。

    如果你想要单位矩阵,只需映射范围:

    val newY = spark.sparkContext.range(0, 8000).map(j => MatrixEntry(j, j, 1))
    

    具有并行化的循环无法扩展,并将所有数据保存在驱动程序内存中Why does SparkContext.parallelize use memory of the driver?

    【讨论】:

      猜你喜欢
      • 2017-11-29
      • 1970-01-01
      • 2017-05-19
      • 2019-11-01
      • 1970-01-01
      • 1970-01-01
      • 2013-12-09
      • 2022-11-13
      • 2016-12-09
      相关资源
      最近更新 更多