如何在Hadoop MapReduce中提取每个标题的最高值对应行?
实现方案:筛选每个title对应value最高的所有记录
当然可以实现这个需求!下面我会给你几种适配Hadoop生态的方案,你可以根据自己的技术栈和数据规模来选择:
方案一:用Hive/Spark SQL(最简便高效)
如果你的数据已经在Hadoop生态里,用SQL类工具是最快的实现方式,不需要写复杂的MapReduce代码。
1. 创建并加载数据(以Hive为例)
首先创建匹配数据格式的表结构:
CREATE TABLE IF NOT EXISTS title_month_value ( title STRING, month STRING, value INT ) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' STORED AS TEXTFILE;
然后把你的Hadoop生成的数据加载到表中(替换成你的实际HDFS数据路径):
LOAD DATA INPATH '/user/hadoop/your_dataset_path' INTO TABLE title_month_value;
2. 执行查询获取结果
使用窗口函数RANK()来保留所有value为最大值的记录(如果用ROW_NUMBER()会只保留一条,不符合你的需求):
SELECT title, month, value FROM ( SELECT title, month, value, -- 按title分组,按value降序排名,相同最大值的记录排名相同 RANK() OVER (PARTITION BY title ORDER BY value DESC) AS rank_num FROM title_month_value ) ranked_data WHERE rank_num = 1;
执行这个查询后,就能得到每个title下所有value为最大值的记录,比如示例里的Fly对应的08、09、11月记录都会被完整保留。
方案二:用MapReduce二次处理(原生MapReduce实现)
如果必须用原生MapReduce来实现,可以分两个阶段完成:
阶段一:计算每个title的最大value值
- Mapper:读取每条
title,month,value记录,输出key=title,value=Integer.valueOf(value) - Reducer:对每个title对应的所有value值取最大值,输出
key=title,value=max_value(比如Fly对应4)
阶段二:关联原数据,筛选符合条件的记录
- Mapper:
- 读取原数据集时,输出
key=title,value="original\t"+month+"\t"+value - 读取阶段一的结果时,输出
key=title,value="max\t"+max_value
- 读取原数据集时,输出
- Reducer:
- 先遍历所有输入值,收集当前title的max_value(每个title只会有一个)
- 再遍历所有原数据记录,判断记录中的value是否等于max_value
- 符合条件的记录,输出
title+","+month+","+value
方案三:用Spark Core(可选)
如果你的集群部署了Spark,用Spark Core实现会比原生MapReduce简洁很多:
import org.apache.spark.SparkContext import org.apache.spark.SparkConf object MaxValueRecords { def main(args: Array[String]) { val conf = new SparkConf().setAppName("MaxValueRecords") val sc = new SparkContext(conf) // 读取原数据并转换格式 val data = sc.textFile("/user/hadoop/your_dataset_path") .map(line => { val parts = line.split(",") (parts(0), (parts(1), parts(2).toInt)) // (title, (month, value)) }) // 计算每个title的最大value val maxValues = data.mapValues(_._2).reduceByKey(Math.max) // 关联原数据,筛选出value等于最大值的记录 val result = data.join(maxValues) .filter { case (_, ((month, value), maxVal)) => value == maxVal } .map { case (title, ((month, value), _)) => s"$title,$month,$value" } // 保存结果到HDFS result.saveAsTextFile("/user/hadoop/output_path") sc.stop() } }
内容的提问来源于stack exchange,提问作者Gideok Seong
相关产品推荐
相关产品推荐

