【发布时间】:2017-08-14 15:07:16
【问题描述】:
我有以下程序计算日志文件中“错误”的计数。最后,它的值会打印在控制台中。当程序在 yarn-client 中运行时,它会在控制台中显示累加器正确的值 509,但是当它在 yarn-cluster 模式下运行时,不会显示这样的值。如何以纱线集群模式打印它?
object ErrorLogsCount{
def main(args:Array[String]){
val sc = new SparkContext();
val logsRDD = sc.textFile(args(0),4)
val errorsAcc = sc.accumulator(0,"Errors Accumulator")
val errorsLogRDD = logsRDD.filter(x => x.contains("ERROR"))
errorsLogRDD.persist()
errorsLogRDD.foreach(x => errorsAcc += 1)
errorsLogRDD.collect()
//printing accumulator
println(errorsAcc.name+" = "+errorsAcc)
//Saving results in HDFS
errorsLogRDD.coalesce(1).saveAsTextFile(args(1))
}
}
尝试在 HDP Sandbox 2.4 (Spark 1.6.0) 中运行
【问题讨论】:
标签: scala apache-spark hadoop-yarn