Spark中继承App特质与main方法实现计数器为何输出不同?
为什么继承App特质和main方法的Spark程序统计缺失记录结果不同?
这是个非常典型的Spark分布式执行模型结合Scala特性导致的问题,我来帮你拆解清楚:
先看现象对应的本质原因
1. 为什么main方法版本输出5 0?
Spark的核心是分布式懒执行,你写的map算子属于转换操作(Transformation),代码不会立即执行,直到遇到count()这种行动操作(Action)才会触发任务提交。
missingRecords是Driver进程中的一个局部变量,当map算子被序列化发送到Executor执行时,这个变量会被复制一份传到Executor节点(或者local模式下的线程)。- Executor上执行
map时,修改的是这个副本的值,Driver端的原始missingRecords完全没被触碰。 - 当
count()执行完成后,Driver输出的还是自己手里那个初始值为0的变量,所以结果是0。
简单说:你以为在改同一个变量,其实Driver和Executor手里的是两个完全独立的副本,Executor的修改不会同步回Driver。
2. 为什么继承App特质的版本输出5 2?
这是Scala App特质的初始化机制加上local模式的巧合导致的:
- Scala的
App特质是基于DelayedInit实现的,你的代码逻辑其实是在类的构造函数中执行的,而不是像main方法那样在独立的主线程入口。 - 当你在
local模式下运行Spark时,Executor其实是在同一个JVM进程的不同线程中运行的,而App特质的构造函数中的变量missingRecords,因为作用域的关系,Executor线程直接引用了Driver端的这个变量(没有走序列化复制的流程)。 - 所以Executor线程对
missingRecords的修改会直接反映到Driver端的变量上,最后能得到正确的2。
⚠️ 注意:这只是local模式下的特例!如果换成集群模式(比如YARN、Standalone),App版本同样会输出0,因为此时Executor在不同的JVM进程,变量还是会被序列化复制,修改的是副本。
正确的Spark统计方式:使用累加器
上面两种写法都不是Spark中统计分布式数据的正确姿势,Spark提供了**累加器(Accumulator)**专门解决这类跨Driver-Executor的计数问题,它能保证修改会正确同步回Driver。
示例代码如下:
object CorrectCounter extends App { val conf = new SparkConf().setAppName("CounterTest").setMaster("local") val sc = new SparkContext(conf) val fileRDD = sc.textFile("data-error.csv") // 创建一个Long类型的累加器,初始值0 val missingRecordsAcc = sc.longAccumulator("MissingRecordsCounter") val rdd1 = fileRDD.map(rec => { val parseResult = RecordParser.parse(rec) if (parseResult.isLeft) { missingRecordsAcc.add(1) // 累加器计数 } rec }) println(rdd1.count()) println(missingRecordsAcc.value) // 获取累加器最终值 }
或者用main方法版本,逻辑完全一样,不管是local还是集群模式,结果都会是5 2。
内容的提问来源于stack exchange,提问作者z_1_p
相关产品推荐
相关产品推荐

