Hadoop MapReduce两阶段任务遭遇FileNotFoundException求助
问题:MapReduce两阶段程序中FileNotFoundException排查与解决
我正在开发一个包含两个阶段的MapReduce程序:第一阶段生成包含两组键值对的文件,用于后续平均值计算;第二阶段基于该平均值执行计算并写入结果文件。目前第一阶段已成功生成文件,但调用计算平均值的方法时触发了FileNotFoundException。
相关代码
import org.apache.commons.text.StringTokenizer; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.*; import java.io.IOException; public class ge3ex3 { private final static String PATH = "avOutput"; public static class ge3ex3AverageMapper extends Mapper<Object, Text, Text, LongWritable>{ //constant variable one private final LongWritable one = new LongWritable(1); //contant variable for cases private final LongWritable cases = new LongWritable(); public void map(Object key, Text value, Context context) throws IOException, InterruptedException { //read the rows String row = value.toString(); //Take the items from each row String[] items = row.split(","); //check if there are 12 columns if (items.length == 12){ try { //read the case column int eachCase = Integer.parseInt(items[4]); //set the count of the cases cases.set(eachCase); //write to context the cases' count and how many records context.write(new Text("cases"), cases); context.write(new Text("records"), one); }catch (NumberFormatException x){ } } } } public static class ge3ex3Reducer extends Reducer<Text, LongWritable, Text, LongWritable>{ private LongWritable sum = new LongWritable(); public void reduce(Text key, Iterable<LongWritable> values, Context context) throws IOException, InterruptedException{ int tempSum = 0; //Iterate through values and add up the sum for(LongWritable val : values){ tempSum += val.get(); } sum.set(tempSum); //write the sum into context context.write(key, sum); } } public static class ge3ex3Mapper extends Mapper<Object, Text, Text, LongWritable>{ private final static LongWritable ONE = new LongWritable(1); public void map(Object key, Text value, Context context) throws IOException, InterruptedException{ //read the file from the first mapper Configuration conf = context.getConfiguration(); Double av = Double.parseDouble(conf.get("av")); //pass each value to string String row = value.toString(); //pass it to array splited by comma String[] items = row.split(","); //check if the row has more cases than the average if(items.length == 12){ try { int eachCase = Integer.parseInt(items[4]); if(eachCase > av){ //take year, month, country and write them into context String year = items[3]; String month = items[2]; String country = items[6]; context.write(new Text(year + " / " + month + " / " + country), ONE); } }catch (NumberFormatException x){ } } } } //calculate the average private static double calculateAv()throws IOException{ //store the sum of cases & records long cases = 0; long records = 0; //read from the file from the first map-reduce phase File fs = new File(PATH + "/part-r-00000"); if (!fs.exists()) throw new FileNotFoundException("File not Found: " + fs.getPath()); //read from file BufferedReader bt = new BufferedReader(new FileReader(fs)); String line; //while the line is not empty while ((line = bt.readLine()) != null){ //Parse each line StringTokenizer sr = new StringTokenizer(line); // String value = sr.nextToken(); //check if value is CASES or RECORDS //and sign them into their perspective variables if(value.equals("cases")){ String countCases = sr.nextToken(); cases = Long.parseLong(countCases); }else if (value.equals("records")){ String countRecords = sr.nextToken(); records = Long.parseLong(countRecords); } } //calculate the average and return it double av = cases / (double) records; return av; } public static void main(String[] args)throws Exception{ Configuration conf = new Configuration(); Job job = Job.getInstance(conf,"ge3ex3_get_Average"); job.setJarByClass(ge3ex3.class); job.setMapperClass(ge3ex3AverageMapper.class); job.setReducerClass(ge3ex3Reducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(LongWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(PATH)); //wait for the next map - reduce phase to complete job.waitForCompletion(true); //calculate the average based on the first mapper cycle double av = calculateAv(); //Second Mapping Phase //pass the list with the average number conf.setDouble("av", av); Job job1 = Job.getInstance(conf, "ge3ex3_result"); job1.setJarByClass(ge3ex3.class); job1.setMapperClass(ge3ex3Mapper.class); job1.setReducerClass(ge3ex3Reducer.class); job1.setOutputKeyClass(Text.class); job1.setOutputValueClass(LongWritable.class); FileInputFormat.addInputPath(job1, new Path(args[0])); FileOutputFormat.setOutputPath(job1, new Path(args[1])); //wait for the next map - reduce phase to complete job1.waitForCompletion(true); } }
错误日志
Exception in thread "main" java.io.FileNotFoundException: File not Found: avOutput/part-r-00000 at ge3ex3.calculateAv(ge3ex3.java:95) at ge3ex3.main(ge3ex3.java:134) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.hadoop.util.RunJar.run(RunJar.java:323) at org.apache.hadoop.util.RunJar.main(RunJar.java:236)
core.xml配置
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> </configuration>
问题原因与解决方案
核心原因
第一阶段的输出文件写入到了HDFS文件系统(从core.xml的fs.defaultFS配置可知),但calculateAv()方法中使用了Java标准的File和FileReader类,这些类只能访问本地文件系统,无法读取HDFS上的文件,因此触发找不到文件的异常。
修复步骤
1. 改用HDFS API读取输出文件
修改calculateAv()方法,使用Hadoop的FileSystem和FSDataInputStream来读取HDFS上的文件:
private static double calculateAv(Configuration conf) throws IOException{ long cases = 0; long records = 0; // 获取HDFS文件系统实例 FileSystem fs = FileSystem.get(conf); Path outputPath = new Path(PATH + "/part-r-00000"); if (!fs.exists(outputPath)) { throw new FileNotFoundException("File not Found: " + outputPath.toString()); } // 使用HDFS API读取文件 try (BufferedReader bt = new BufferedReader(new InputStreamReader(fs.open(outputPath)))) { String line; while ((line = bt.readLine()) != null) { StringTokenizer sr = new StringTokenizer(line); String key = sr.nextToken(); if(key.equals("cases")){ String countCases = sr.nextToken(); cases = Long.parseLong(countCases); }else if (key.equals("records")){ String countRecords = sr.nextToken(); records = Long.parseLong(countRecords); } } } return cases / (double) records; }
2. 修改main方法调用
在main方法中调用calculateAv()时,传入Configuration对象:
// 原代码:double av = calculateAv(); double av = calculateAv(conf);
3. 额外注意事项
- 确保HDFS服务正常运行,可通过
hdfs dfs -ls avOutput命令验证文件是否存在 - 若运行时出现权限问题,可给输出目录添加权限:
hdfs dfs -chmod 755 avOutput - 第一阶段的输出路径
PATH如果是相对路径,会基于HDFS的用户根目录(如/user/yourname/avOutput),可通过hdfs dfs -ls查看确认
内容的提问来源于stack exchange,提问作者Twofaced Tough
相关产品推荐
相关产品推荐

