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