共享变量之广播变量
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))
广播变量的内部机制
分发过程:
- Driver 把广播变量切成小块(默认 4MB/块)
- 块元数据(块在哪)广播给所有 Executor
- Executor 从 Driver 或其他 Executor 拉取块数据(类似 BT 下载)
- 拼装完整数据,缓存在 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 只写 |
| 典型用途 | 共享大字典、配置 | 计数、求和、收集错误 |
| 类比 | 每个工人发一本操作手册 | 每个工人汇报工作量 |