习题 9:广播变量 join 优化
题目:有小表 smallTable(公司信息)和大表
bigTable(订单信息),使用广播变量实现高效 join。
import org.apache.spark.{SparkConf, SparkContext}
object CsvJoinExample {
def main(args: Array[String]): Unit = {
val conf = new SparkConf().setAppName("CsvJoin").setMaster("local[*]")
val sc = new SparkContext(conf)
// 1. 读取小表(价格映射)并转为 Map、广播
val smallTableRaw = sc.textFile("prices.csv")
.filter(!_.startsWith("clusterName")) // 跳过表头
.map { line =>
val parts = line.split(",")
(parts(0), parts(1).toInt) // (clusterName, price)
}
.collectAsMap()
val smallTableBroadcast = sc.broadcast(smallTableRaw)
// 2. 读取大表(订单数据)
val bigTable = sc.textFile("orders.csv")
.filter(!_.startsWith("orderId")) // 跳过表头
.map { line =>
val parts = line.split(",")
(parts(0), parts(1), parts(2)) // (orderId, clusterName, year)
}
// 3. 使用广播变量进行 map-side join
val result = bigTable.map { case (orderId, clusterName, year) =>
val price = smallTableBroadcast.value.getOrElse(clusterName, "unknown")
(orderId, clusterName, price, year)
}
// 4. 输出结果
result.collect().foreach(println)
sc.stop()
}
}