习题 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()
  }
}