【发布时间】:2015-12-08 22:40:29
【问题描述】:
我正在尝试在集群任务中访问累加器的值。但是当我这样做时会引发异常:
无法读取累加器的值
我尝试使用row.localValue,但它返回相同的数字。有解决办法吗?
private def modifyDataset(
data: String, row: org.apache.spark.Accumulator[Int]): Array[Int] = {
var line = data.split(",")
var lineSize = line.size
var pairArray = new Array[Int](lineSize-1)
var a = row.value
paiArray(0)=a
row+=1
pairArray
}
var sc = Spark_Context.InitializeSpark
var row = sc.accumulator(1, "Rows")
var dataset = sc.textFile("path")
var pairInfoFile = noHeaderRdd.flatMap{ data => modifyDataset(data,row) }
.persist(StorageLevel.MEMORY_AND_DISK)
pairInfoFile.count()
【问题讨论】:
标签: scala apache-spark accumulator