【问题标题】:How to access the value of accumulator in tasks?如何访问任务中累加器的值?
【发布时间】: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


    【解决方案1】:

    这根本不可能,也没有解决方法。 Spark accumulators 从工作人员的角度来看是只写变量。在任务期间读取其值的任何尝试都是没有意义的,因为工作人员之间没有共享状态,并且本地累加器值仅反映当前分区的状态。

    一般而言,accumulators 主要用于诊断,不应用作应用程序逻辑的一部分。在转换中使用时,您获得的唯一保证是至少执行一次。

    另请参阅:How to print accumulator variable from within task (seem to "work" without calling value method)?

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-10-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-02-15
      • 1970-01-01
      相关资源
      最近更新 更多