Spark LongAccumulator终值始终为0问题排查求助
Spark Streaming中LongAccumulator最终值始终为0的问题分析与解决
这个问题我太熟悉了,本质是没搞清楚Spark Streaming的惰性执行模型和累加器的更新时机!
问题根源
你现在的代码执行顺序是这样的:
- 在
startJob方法里,你定义了DStream的转换(map里累加器add)和输出操作(saveToCassandra),但这些都只是定义任务逻辑,并不会立刻执行——Spark Streaming的所有DStream操作都是惰性的,只有当StreamingContext启动后,每个批次到来时才会真正调度任务到Executor执行。 - 你在调用
saveToCassandra之后立刻打印accTotal.value,这时候StreamingContext还没启动,没有任何任务运行,累加器自然还是初始值0。 - 之后你才启动
StreamingContext,这时候批次任务开始执行,Executor里的map操作才会真正更新累加器——这时候你之前的打印语句早就执行完了,所以看不到更新后的值。
而你在map里看到的累加器递增,以及UI里的更新,都是任务执行时Executor端的实时状态,但Driver主线程的打印是在任务启动前触发的,完全不在一个时间点上。
解决方案
要拿到更新后的累加器值,必须等到对应的批次任务执行完成后再获取,推荐两种方式:
方式一:用foreachRDD在批次处理后打印
foreachRDD是DStream提供的、能在Driver端访问每个批次RDD执行结果的方法,你可以把累加器的打印逻辑放在这里,确保是在当前批次任务执行后获取值:
def saveToCassandra(upserts: DStream[Data]) = { upserts.foreachRDD { rdd => val countedRDD = rdd.map { record => accTotal.add(1) println("ACC: " + accTotal.value) // Executor端实时打印 DataExt(record.data, record.data1) } countedRDD.saveToCassandra(keyspace, table) // 这里是Driver端,当前批次任务执行后打印累加器值 println("XXX:" + accTotal.value) } }
方式二:通过StreamingListener监听批次完成事件
如果需要更全局的监控,可以自定义一个监听器,在每个批次完成后自动获取累加器值:
// 自定义Streaming监听器 class AccumulatorMonitor(acc: LongAccumulator) extends StreamingListener { override def onBatchCompleted(event: StreamingListenerBatchCompleted): Unit = { println(s"批次 ${event.batchInfo.batchTime} 完成,累计写入记录数:${acc.value}") } } // 在startJob里注册监听器 override def startJob(ssc: StreamingContext): Unit = { accTotal = ssc.sparkContext.longAccumulator("total") // 注册监听器 ssc.addStreamingListener(new AccumulatorMonitor(accTotal)) val inputKafka = createDirectStream(ssc, kafkaParams, topicsSet) val rddAvro = inputKafka.map{x => x.value()} saveToCassandra(rddAvro) // 这里不再直接打印,交给监听器处理 }
注意事项
- 永远不要在
StreamingContext启动前的Driver主线程里直接打印累加器值——这时候任务还没执行,累加器肯定是初始值。 - Spark累加器的更新是Executor端任务执行时触发,然后异步同步到Driver的,所以必须等待任务完成后才能拿到最新值。
内容的提问来源于stack exchange,提问作者Guille
相关产品推荐
相关产品推荐

