仅被单个任务使用时Spark Broadcast Variable的使用优势问询
Spark单场景使用广播变量的额外优势
即使明确变量仅在作业中使用一次,广播变量对比普通闭包变量依然存在以下不可忽略的优势:
- 传输效率优势:普通闭包变量是跟随每个Task序列化分发,同一Executor上运行多少个关联Task,就会重复传输多少次变量;广播变量仅向每个Executor分发一次副本,与该节点上的Task数量无关。哪怕作业只有一次Action,只要并行度高于Executor数量,就能减少传输量,变量体积越大优势越明显。
- 内存占用优势:普通闭包变量会在每个Task执行时单独反序列化出独立副本,同一Executor上的多份副本会占用多倍内存,增加GC压力;广播变量在每个Executor上仅保留一份副本,所有Task共享使用,可大幅降低内存冗余。
- 调度性能优势:普通变量会增大每个Task的序列化体积,提升Driver端序列化Task、Executor端反序列化Task的开销,高并行度场景下会明显拖慢调度流程;广播变量的传输独立于Task调度流程,不会增加Task本身的调度开销。
仅在变量体积极小(KB级别)、作业并行度极低的场景下,广播变量的协调开销会高于普通闭包变量,此时可选择使用普通变量。
对应的Java示例代码如下:
public class SparkDriver { public static void main(String[] args) { String inputPath = args[0]; String outputPath = args[1]; Map<String,String> dictionary = new HashMap<>(); dictionary.put("J", "Java"); dictionary.put("S", "Spark"); SparkConf conf = new SparkConf() .setAppName("广播变量测试") .setMaster("local"); try (JavaSparkContext context = new JavaSparkContext(conf)) { // 注册字典为广播变量 final Broadcast<Map<String,String>> dictionaryBroadcast = context.broadcast(dictionary); context.textFile(inputPath) .map(line -> { // 仅在该转换算子中使用一次广播变量 Map<String,String> d = dictionaryBroadcast.value(); String[] words = line.split(" "); StringBuffer sb = new StringBuffer(); for (String w : words) sb.append(d.get(w)).append(" "); return sb.toString(); }) .saveAsTextFile(outputPath); // 作业仅执行一次Action算子 } } }
内容的提问来源于stack exchange,提问作者Franco Ruggeri
相关产品推荐
相关产品推荐

