Spark广播变量能否变更?示例中值改变的原因解析
Spark广播变量"可变"的真相解答
先来看你贴的这段代码:
public class TestBroadcast { private static JavaStreamingContext jssc; private static volatile Broadcast<Boolean> done; public static void main(String[] args) throws InterruptedException { Logger.getLogger("org").setLevel(Level.ERROR); List<String> log = Arrays.asList("X", "X", "X"); SparkConf sparkConf = new SparkConf().setAppName("Test").setMaster("local[2]"); jssc = new JavaStreamingContext(sparkConf, Durations.seconds(1)); done = jssc.sparkContext().broadcast(Boolean.FALSE); // false in the beginning JavaRDD<String> _rdd = jssc.sparkContext().parallelize(log); Queue<JavaRDD<String>> queue = new LinkedList<>(); queue.add(_rdd); JavaDStream<String> lines = jssc.queueStream(queue); lines.foreachRDD( rdd -> { rdd.foreachPartition(x -> System.out.println(done.getValue())); // executor get false done = jssc.sparkContext().broadcast(Boolean.TRUE); // driver set the variable to true /*MARK*/ rdd.foreachPartition(x -> System.out.println(done.getValue())); // executor get true }); jssc.start(); jssc.awaitTermination(); } }
核心问题解答:广播变量是否有可能发生变更?
严格来说,Spark原生的广播变量实例是不可变的——一旦你通过SparkContext.broadcast()创建了一个广播变量,它内部封装的值就被固定了,Executor端拿到的是这个值的只读副本,根本没法修改它。
但你看到的代码并不是修改了原广播变量,而是在Driver端做了个"偷梁换柱"的操作:重新创建了一个新的广播变量,然后把原来的done引用(这个是Driver端的变量)指向了新的广播对象。这本质是替换了引用,不是修改原有广播变量的内容。
何时会出现"变量变更"的假象?
只有当Driver端主动重新创建广播变量并覆盖原有引用的时候,才会让后续的任务拿到新的值:
- 旧的广播变量其实还存在,要是有Executor任务还持有旧的广播变量引用,拿到的依然是原来的值;
- 新创建的广播变量会被Spark重新分发给Executor,之后提交的任务使用更新后的引用,就能拿到新值。
为什么标注/MARK/的代码行之后,值发生了改变?
我们一步步拆解这段代码的执行逻辑:
- 程序启动时,
done指向第一个广播变量,值是false。第一个foreachPartition任务被提交到Executor,此时用的是旧的done引用,所以Executor拉取的是旧广播变量,打印false。 - 回到Driver端,执行
done = jssc.sparkContext().broadcast(Boolean.TRUE);,这时候done这个变量已经指向了一个全新的、值为true的广播变量。 - 第二个
foreachPartition任务是在Driver更新了done引用之后才提交的,所以Executor会拉取这个新的广播变量,自然就打印出true了。
这里要注意:这两个foreachPartition是在同一个foreachRDD的Driver逻辑里先后提交的任务,第二个任务用的是已经更新后的引用,所以能拿到新值。
内容的提问来源于stack exchange,提问作者petertc
相关产品推荐
相关产品推荐

