Spark技术问询:如何按年份输出数据集的前10行
嘿,我来帮你搞定这个需求!看起来你已经完成了词频统计的部分,现在要按年份拆分并输出每个年份的前10条数据对吧?咱们一步步来调整和补充代码:
第一步:从现有RDD中提取年份
你的counts RDD里,键是类似2004-dog这种带年份的字符串,值是对应的计数值。首先咱们要把年份从键里提取出来,把数据转换成(年份, (原键, 计数值))的结构,方便后续分组:
// 从counts的键中提取前4位作为年份,转换为(年份, (原键, 计数值))的PairRDD JavaPairRDD<String, Tuple2<String, Integer>> yearKeyValuePairs = counts.mapToPair(tuple -> { String originalKey = tuple._1(); // 提取前4个字符作为年份 String year = originalKey.substring(0, 4); return new Tuple2<>(year, new Tuple2<>(originalKey, tuple._2())); });
第二步:按年份分组并取前10条
这里推荐用aggregateByKey来处理,比直接groupByKey更高效——它会先在每个分区做局部聚合,再合并结果,避免把大量数据拉到单个节点。
情况1:按计数值降序取前10条(热门词)
如果想每个年份下取计数最高的前10个单词,用下面的代码:
// 初始化每个年份的空列表 Function0<List<Tuple2<String, Integer>>> createCombiner = () -> new ArrayList<>(); // 单个分区内,把新元素加入列表,排序后保留前10 Function2<List<Tuple2<String, Integer>>, Tuple2<String, Integer>, List<Tuple2<String, Integer>>> mergeValue = (list, item) -> { list.add(item); // 按计数值降序排序 list.sort((a, b) -> Integer.compare(b._2(), a._2())); // 只保留前10条 return list.size() > 10 ? list.subList(0, 10) : list; }; // 合并两个分区的列表,同样排序后保留前10 Function2<List<Tuple2<String, Integer>>, List<Tuple2<String, Integer>>, List<Tuple2<String, Integer>>> mergeCombiners = (list1, list2) -> { list1.addAll(list2); list1.sort((a, b) -> Integer.compare(b._2(), a._2())); return list1.size() > 10 ? list1.subList(0, 10) : list1; }; // 执行aggregateByKey,得到每个年份的前10条热门数据 JavaPairRDD<String, List<Tuple2<String, Integer>>> top10ByYear = yearKeyValuePairs.aggregateByKey( createCombiner, mergeValue, mergeCombiners );
情况2:保留原始顺序取前10条
如果不需要排序,只要每个年份下的前10条原始数据,去掉排序逻辑即可:
Function0<List<Tuple2<String, Integer>>> createCombiner = () -> new ArrayList<>(); Function2<List<Tuple2<String, Integer>>, Tuple2<String, Integer>, List<Tuple2<String, Integer>>> mergeValue = (list, item) -> { if (list.size() < 10) { list.add(item); } return list; }; Function2<List<Tuple2<String, Integer>>, List<Tuple2<String, Integer>>, List<Tuple2<String, Integer>>> mergeCombiners = (list1, list2) -> { List<Tuple2<String, Integer>> combined = new ArrayList<>(); combined.addAll(list1); // 从第二个列表补充,直到满10条 for (Tuple2<String, Integer> item : list2) { if (combined.size() < 10) { combined.add(item); } else { break; } } return combined; }; JavaPairRDD<String, List<Tuple2<String, Integer>>> top10ByYear = yearKeyValuePairs.aggregateByKey( createCombiner, mergeValue, mergeCombiners );
第三步:打印结果
最后遍历top10ByYear,把每个年份的前10条数据打印出来:
top10ByYear.foreach(tuple -> { String year = tuple._1(); List<Tuple2<String, Integer>> top10Items = tuple._2(); System.out.println("=== 年份 " + year + " 的前10条数据 ==="); for (int i = 0; i < top10Items.size(); i++) { Tuple2<String, Integer> item = top10Items.get(i); System.out.printf("%d. %s: %d%n", i+1, item._1(), item._2()); } System.out.println(); });
注意事项
- 记得导入所需的类:
org.apache.spark.api.java.function.Function0、org.apache.spark.api.java.function.Function2、scala.Tuple2、java.util.ArrayList、java.util.List等。 - 如果你的原始数据集不是经过词频统计的,而是直接每行是
"2004-dog" 45这种结构,那可以跳过你已有的mentions和counts代码,直接从原始数据提取年份分组即可。
内容的提问来源于stack exchange,提问作者Danny
相关产品推荐
相关产品推荐

