如何在MapReduce中实现获取播放量Top10的电影及名称?
嘿,我来帮你搞定这个Top10的需求!你现在已经完成了电影数据和评分数据的关联,成功统计出了每部电影的播放次数,接下来只需要在现有代码基础上添加Top10的筛选逻辑就行。这里给你两种实用方案,按需选择:
方案一:单个Reducer维护全局Top10(推荐,简单直接)
这种方法适合数据量不是特别夸张的场景,核心思路是在Reducer里用最小堆来维护当前播放量最高的10部电影——堆顶始终是当前Top10里播放量最小的那个,每次新电影的播放次数进来,如果比堆顶大,就替换堆顶,这样堆里始终保留着最大的10个值。最后在Reducer的cleanup阶段把堆里的元素倒序输出就是Top10了。
修改后的完整Reducer代码
import java.io.IOException; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; import java.util.PriorityQueue; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class MoviesRatingJoinReducer extends Reducer<Text, Text, Text, Text> { private ArrayList<Text> listMovies = new ArrayList<Text>(); private ArrayList<Text> listRating = new ArrayList<Text>(); // 用最小堆维护Top10电影 private PriorityQueue<MovieCount> top10Queue; // 内部类存储电影名和播放次数 static class MovieCount { int count; String title; public MovieCount(int count, String title) { this.count = count; this.title = title; } } @Override protected void setup(Context context) throws IOException, InterruptedException { // 初始化最小堆:堆顶是当前Top10中播放量最小的元素 top10Queue = new PriorityQueue<>(10, new Comparator<MovieCount>() { @Override public int compare(MovieCount o1, MovieCount o2) { return Integer.compare(o1.count, o2.count); } }); } @Override public void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { listMovies.clear(); listRating.clear(); for (Text text : values) { if (text.charAt(0) == 'M') { listMovies.add(new Text(text.toString().substring(1))); } else if (text.charAt(0) == 'R') { listRating.add(new Text(text.toString().substring(1))); } } executeJoinLogic(context); } private void executeJoinLogic(Context context) throws IOException, InterruptedException { if (!listMovies.isEmpty() && !listRating.isEmpty()) { int playCount = listRating.size(); for (Text moviesData : listMovies) { String movieTitle = moviesData.toString(); // 维护Top10堆 if (top10Queue.size() < 10) { top10Queue.add(new MovieCount(playCount, movieTitle)); } else { // 如果当前电影播放量比堆顶大,替换堆顶 if (playCount > top10Queue.peek().count) { top10Queue.poll(); top10Queue.add(new MovieCount(playCount, movieTitle)); } } } } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { // 把堆转成列表并按播放量降序排序 ArrayList<MovieCount> top10List = new ArrayList<>(top10Queue); Collections.sort(top10List, new Comparator<MovieCount>() { @Override public int compare(MovieCount o1, MovieCount o2) { // 倒序排列,让播放量高的在前 return Integer.compare(o2.count, o1.count); } }); // 输出Top10结果 for (MovieCount mc : top10List) { context.write(new Text(mc.title), new Text(String.valueOf(mc.count))); } } }
关键注意点
- 记得在你的MapReduce Driver类里设置Reducer数量为1:
job.setNumReduceTasks(1);,否则多个Reducer会各自统计分片内的Top10,得不到全局结果。 - 最小堆的优势是内存占用固定(最多10个元素),不会因为数据量太大导致内存溢出。
方案二:二次MapReduce实现全局Top10(适合超大数据量)
如果你的数据集特别大,单个Reducer处理不过来,可以分两步走:
- 第一步:用你现有的作业完成电影播放次数统计,输出格式为
电影名::播放次数。 - 第二步:写一个新的MapReduce作业:
- Mapper:读取第一步的输出,把播放次数转成
IntWritable作为Key,电影名作为Value(注意要把Key设为倒序排序,这样Reducer会先收到播放量最高的记录)。 - Reducer:只取前10条记录输出,直接得到全局Top10。
- Mapper:读取第一步的输出,把播放次数转成
这种方案扩展性更好,但需要写两个作业,稍微复杂一点。
内容的提问来源于stack exchange,提问作者Jeet
相关产品推荐
相关产品推荐

