共享变量之广播变量

Spark 广播变量:大字典表别让每个 Task 都传一份,发给每个 Executor 一份就够了

Spark 里每个 Task 执行时,闭包里的变量会被复制一份到该 Task。如果变量很小(几 KB),没问题。但如果变量很大(比如 1GB 的字典表),100 个 Task 就是 100 份拷贝。

分发方式 拷贝数 网络传输 内存占用
普通闭包 每个 Task 一份 大(100 × 1GB) 大(100 × 1GB)
广播变量 每个 Executor 一份 小(Executor 数 × 1GB) 小(Executor 数 × 1GB)

广播变量就是:把大对象从 Driver 发给每个 Executor,只发一次,该 Executor 上的所有 Task 共享这一份。

广播变量的基本用法

// 1. 创建广播变量
val dict = Map("a" -> 10, "b" -> 20, "c" -> 30)
val bcDict = sc.broadcast(dict)

// 2. Task 里用 broadcast.value 取数据
val rdd = sc.parallelize(List("a", "b", "c", "d"))
val result = rdd.map(key => bcDict.value.getOrElse(key, 0))

// 3. 不需要时释放(可选)
bcDict.unpersist()

关键点:

  • 广播变量是只读的——Executor 上只能读,不能改
  • 广播变量的值在 Executor 上只加载一次,所有 Task 共享
  • 使用 .value 访问广播的内容

广播变量 vs 普通闭包

//  普通闭包:每个 Task 都传一份 dict
val dict = Map("a" -> 10, "b" -> 20)   // 假设 1GB
val rdd = sc.parallelize(1 to 1000, 100)  // 100 个 Task
rdd.map(key => dict.getOrElse(key, 0)).collect()
// 网络传输:100 × 1GB = 100GB

// 广播变量:每个 Executor 只传一份
val bcDict = sc.broadcast(dict)
rdd.map(key => bcDict.value.getOrElse(key, 0)).collect()
// 网络传输:Executor 数 × 1GB(假设 10 个 Executor = 10GB)

典型使用场景

1. Map 端 Join(小表广播)

小表 Join 大表时,把小表广播到每个 Executor,在 Map 端完成关联,不走 Shuffle。

// DataFrame 方式(推荐)
import org.apache.spark.sql.functions.broadcast

val smallDF = spark.read.parquet("small_table")
val largeDF = spark.read.parquet("large_table")

val result = largeDF.join(broadcast(smallDF), "key")

2. 字典映射

// 把字典表广播出去,每条数据过来查一下
val dict = Map("product_001" -> "电子产品", "product_002" -> "服装")
val bcDict = sc.broadcast(dict)

rdd.map(record => {
  val productId = record.productId
  val category = bcDict.value.getOrElse(productId, "未知")
  (record, category)
})

3. 共享配置参数

// 模型参数、配置信息广播出去,所有 Task 用同一套
val config = Config(learningRate = 0.01, batchSize = 128)
val bcConfig = sc.broadcast(config)

rdd.map(data => train(data, bcConfig.value))

广播变量的内部机制

分发过程:

  1. Driver 把广播变量切成小块(默认 4MB/块)
  2. 块元数据(块在哪)广播给所有 Executor
  3. Executor 从 Driver 或其他 Executor 拉取块数据(类似 BT 下载)
  4. 拼装完整数据,缓存在 Executor 内存里

优势: Driver 不是唯一数据源,多个 Executor 互相拉取,避免 Driver 网络成为瓶颈。

使用注意事项

1. 广播变量不能太大

Executor 内存有限,广播变量太大(比如 > 2GB)会 OOM。

如果数据太大,考虑:

  • 用分布式存储(HDFS)共享,不用广播
  • 分块广播,分批处理

2. 广播变量是只读的

广播变量的值不可修改。如果要在 Executor 上修改数据,用累加器(Accumulator)。

3. 及时释放

用完后 unpersist(),释放 Executor 内存。

bcDict.unpersist()

4. 序列化问题

广播变量里的对象需要可序列化(实现 Serializable)。用 Kryo 序列化可以更高效:

conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")

广播变量 vs 累加器

广播变量 累加器
数据流向 Driver → Executor(分发) Executor → Driver(聚合)
权限 Executor 只读 Executor 只写
典型用途 共享大字典、配置 计数、求和、收集错误
类比 每个工人发一本操作手册 每个工人汇报工作量