Hadoop MapReduce实现访问量升序排序:键值交换逻辑与数据类型适配问题
问题分析与解决方案
看起来你想通过交换键值对让Hadoop自动排序访问量,但你的代码在逻辑和数据类型匹配上都出现了问题,导致无法得到期望的聚合排序结果。我帮你梳理下问题点,然后给出修正后的实现:
核心问题拆解
- Mapper逻辑错误:你现在的Mapper直接把每条日志的访问量(固定为1)作为key,URL作为value输出,这会导致Shuffle阶段把所有key=1的条目归到同一个Reducer,根本无法实现按URL聚合访问量的第一步。正确的第一步应该是先按URL分组统计访问次数。
- Reducer数据类型不匹配:你的Reducer定义的输入键值是
IntWritable, Text,但reduce方法的参数却写了Text key, Iterable<IntWritable> values,这完全不匹配,会直接导致编译错误,这也是你遇到数据类型问题的根源。 - 排序逻辑缺失:要实现按访问量升序,需要先完成URL访问量的统计,再将统计结果的键值交换(访问量作为key,URL作为value),利用Hadoop Shuffle阶段的排序特性,最后再还原成URL+访问量的格式。
修正后的实现方案
我们需要分两个Job完成:第一个Job统计每个URL的访问量,第二个Job对统计结果按访问量升序排序。
第一个Job:统计URL访问量
Mapper类(统计阶段)
import java.io.IOException; import java.util.regex.Matcher; import java.util.regex.Pattern; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class URLCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text url = new Text(); private Pattern httplogPattern = Pattern.compile("([^\\s]+) - - \\[(.+)\\] \"([^\\s]+) (/[^\\s]*) HTTP/[^\\s]+\" [^\\s]+ ([0-9]+)"); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); Matcher matcher = httplogPattern.matcher(line); if (matcher.matches()) { url.set(matcher.group(1)); // 提取URL作为key context.write(url, one); // 输出<URL, 1> } } }
Reducer类(统计阶段)
import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class URLCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } context.write(key, new IntWritable(sum)); // 输出<URL, 总访问量> } }
第二个Job:按访问量升序排序
Mapper类(排序阶段)
这个Mapper的作用是交换统计结果的键值对,让访问量作为key,触发Hadoop的自动排序:
import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class SortMapper extends Mapper<LongWritable, Text, IntWritable, Text> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 读取第一个Job的输出,格式是"URL 访问量" String[] parts = value.toString().split("\\s+"); if (parts.length == 2) { String url = parts[0]; int count = Integer.parseInt(parts[1]); context.write(new IntWritable(count), new Text(url)); // 输出<访问量, URL> } } }
Reducer类(排序阶段)
这个Reducer的作用是把键值对还原成期望的格式输出:
import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class SortReducer extends Reducer<IntWritable, Text, Text, IntWritable> { @Override protected void reduce(IntWritable key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 同一个访问量可能对应多个URL,逐个输出 for (Text url : values) { context.write(url, key); // 输出<URL, 访问量> } } }
主类(串联两个Job)
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class URLSortDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); // 第一个Job:统计URL访问量 Job countJob = Job.getInstance(conf, "URL访问量统计"); countJob.setJarByClass(URLSortDriver.class); countJob.setMapperClass(URLCountMapper.class); countJob.setReducerClass(URLCountReducer.class); countJob.setOutputKeyClass(Text.class); countJob.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(countJob, new Path(args[0])); Path tempOutput = new Path("temp_output"); FileOutputFormat.setOutputPath(countJob, tempOutput); // 删除临时输出目录(如果存在) FileSystem fs = FileSystem.get(conf); if (fs.exists(tempOutput)) { fs.delete(tempOutput, true); } // 等待第一个Job完成,再执行第二个Job if (countJob.waitForCompletion(true)) { // 第二个Job:按访问量升序排序 Job sortJob = Job.getInstance(conf, "访问量升序排序"); sortJob.setJarByClass(URLSortDriver.class); sortJob.setMapperClass(SortMapper.class); sortJob.setReducerClass(SortReducer.class); sortJob.setOutputKeyClass(IntWritable.class); sortJob.setOutputValueClass(Text.class); FileInputFormat.addInputPath(sortJob, tempOutput); FileOutputFormat.setOutputPath(sortJob, new Path(args[1])); // 删除最终输出目录(如果存在) if (fs.exists(new Path(args[1]))) { fs.delete(new Path(args[1]), true); } System.exit(sortJob.waitForCompletion(true) ? 0 : 1); } else { System.exit(1); } } }
为什么这样能解决问题?
- 第一个Job完成了URL访问量的统计,得到每个URL对应的总访问次数,这是排序的基础。
- 第二个Job通过交换键值对,让访问量作为Mapper的输出key,Hadoop的Shuffle阶段会自动按照key(访问量)进行升序排序,最后在Reducer里还原成<URL, 访问量>的格式输出,就得到了你期望的结果。
- 所有数据类型都严格匹配,避免了之前的类型错误问题。
内容的提问来源于stack exchange,提问作者scandll
相关产品推荐
相关产品推荐

