习题 10:电影评分

题目:读取评分文件和电影文件,计算每部电影的平均评分,筛选平均分大于等于 4.0 的电影,输出电影名和评分。

movies.csv

movieId,title,genres
1,Toy Story (1995),Adventure|Animation|Children|Comedy|Fantasy
2,Jumanji (1995),Adventure|Children|Fantasy

ratings.csv

userId,movieId,rating,timestamp
14,2,4.4,957934979
15,4,1.2,884032602
15,11,4.5,902412488
import org.apache.spark.{SparkConf, SparkContext}

object MovieRatingFixed {
  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setAppName("MovieRating").setMaster("local[*]")
    val sc = new SparkContext(conf)

    val rddrating = sc.textFile(args(0)) // ratings.csv
    val ratingres = rddrating.map(line => line.split(","))
      .filter(a => a(1) != "movieId")
      .map(a => (a(1).toInt, a(2).toDouble))
      .groupByKey()
      .mapValues(a => a.sum / a.size)
      .filter(kv => kv._2 >= 4.0)

    val rddmovies = sc.textFile(args(1)) // movies.csv
      .map(line => line.split(","))
      .filter(a => a(0) != "movieId")
      .map(a => (a(0).toInt, a(1)))
      .join(ratingres)
      .map(f => f._2)

    rddmovies.saveAsTextFile(args(2))
  }
}