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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 12:02:34