You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 04:04:08