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

MapReduce作业输出为空,Reducer无打印输出问题求助

MapReduce作业输出为空的问题排查与解决

你的MapReduce作业输出为空,核心原因是Mapper输出的Value类型不兼容Hadoop的序列化机制,同时存在Job配置缺失、Combiner使用错误等问题,导致Reducer无法接收到Mapper传输的数据。以下是具体问题分析和修复方案:

核心问题分析

1. Mapper输出的Value类型不支持Hadoop序列化

Hadoop要求Map和Reduce之间传输的键值对类型必须实现Writable接口(或其子接口),而你使用的ArrayList<Object>并非Hadoop的序列化类型,Map阶段输出的数据无法被正确序列化传输到Reduce阶段,导致Reducer完全收不到数据,自然没有输出。

2. Job配置缺失Mapper输出类型声明

在Main类中,你只配置了最终作业的输出类型(IntWritable、IntWritable),但没有声明Mapper的输出类型(Text和自定义Value类型)。Hadoop无法自动推断非标准Writable类型的序列化方式,进一步加剧了数据传输失败。

3. Combiner使用错误

你将Reducer3同时设置为Combiner和Reducer,但Combiner的输出类型必须与Mapper的输出类型一致(作为Reducer的输入),而当前Reducer3的输出是IntWritable和IntWritable,与Mapper的输出类型完全不匹配,导致Combiner阶段的数据无法传递到Reducer,甚至直接丢失数据。

4. Arbre类的静态变量线程安全问题

Arbre类中的champs是静态数组,在Map任务的多线程执行环境中,多个线程会同时修改这个静态变量,导致部分记录解析错误,而你的catch块没有任何日志输出,无法发现这些错误。

修复方案

步骤1:自定义Writable类替代ArrayList

创建实现Writable接口的自定义类,用于传输Map阶段的结果:

package TP6;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Writable;

import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;

public class ArbreStatsWritable implements Writable {
    private IntWritable arrondissement;
    private IntWritable age;

    public ArbreStatsWritable() {
        this.arrondissement = new IntWritable();
        this.age = new IntWritable();
    }

    public ArbreStatsWritable(int arrondissement, int age) {
        this.arrondissement = new IntWritable(arrondissement);
        this.age = new IntWritable(age);
    }

    public int getArrondissement() {
        return arrondissement.get();
    }

    public int getAge() {
        return age.get();
    }

    @Override
    public void write(DataOutput out) throws IOException {
        arrondissement.write(out);
        age.write(out);
    }

    @Override
    public void readFields(DataInput in) throws IOException {
        arrondissement.readFields(in);
        age.readFields(in);
    }
}

步骤2:修改Mapper类,使用自定义Writable作为输出Value

package TP6;

import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

import java.io.IOException;

public class Mapper3 extends Mapper<LongWritable, Text, Text, ArbreStatsWritable> {
    @Override
    public void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
        String val = value.toString();
        String[] champs = val.split(";");
        if (champs.length < 6) {
            System.err.println("Invalid line: " + val);
            return;
        }

        try {
            String arrondissementStr = champs[1];
            String anneStr = champs[5];
            int arrondissement = Integer.parseInt(arrondissementStr);
            int age = 2023 - Integer.parseInt(anneStr);
            
            context.write(new Text(arrondissementStr), new ArbreStatsWritable(arrondissement, age));
        } catch (Exception e) {
            System.err.println("Error processing line: " + val);
            e.printStackTrace();
        }
    }
}

注:直接在Mapper中解析数据,避免了Arbre类静态变量的线程安全问题,你也可以修改Arbre类,将champs改为实例变量,每次调用fromLine时初始化。

步骤3:修改Reducer类,适配自定义Writable输入

package TP6;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;

import java.io.IOException;

public class Reducer3 extends Reducer<Text, ArbreStatsWritable, IntWritable, IntWritable> {
    private int maxAge = 0;
    private int keymax = 0;
    private final IntWritable cle = new IntWritable();
    private final IntWritable valeur = new IntWritable();

    @Override
    public void reduce(Text key, Iterable<ArbreStatsWritable> values, Context context) throws IOException, InterruptedException {
        System.out.println("Reducer3: key = " + key);
        
        int localMaxAge = 0;
        int localKeyMax = Integer.parseInt(key.toString());

        for (ArbreStatsWritable value : values) {
            int age = value.getAge();
            if (age > localMaxAge) {
                localMaxAge = age;
                localKeyMax = value.getArrondissement();
            }
        }

        if (localMaxAge > maxAge) {
            maxAge = localMaxAge;
            keymax = localKeyMax;
            cle.set(keymax);
            valeur.set(maxAge);
        }
    }

    @Override
    protected void cleanup(Context context) throws IOException, InterruptedException {
        context.write(cle, valeur);
    }
}

注:修复了原代码中maxAge作为全局变量的问题,改为先计算每个arrondissement的局部最大值,再更新全局最大值,最后在cleanup阶段输出最终结果。

步骤4:修改Main类,完善Job配置

package TP6;

import org.apache.hadoop.conf.Configuration;
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 Main {
    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        Job job = Job.getInstance(conf, "arbre old");
        job.setJarByClass(Main.class);
        
        job.setMapperClass(Mapper3.class);
        job.setReducerClass(Reducer3.class);
        
        job.setMapOutputKeyClass(Text.class);
        job.setMapOutputValueClass(ArbreStatsWritable.class);
        
        job.setOutputKeyClass(IntWritable.class);
        job.setOutputValueClass(IntWritable.class);
        
        FileInputFormat.addInputPath(job, new Path(args[0]));
        FileOutputFormat.setOutputPath(job, new Path(args[1]));
        
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }
}

注:移除了原有的Combiner配置,如需使用Combiner,需单独编写符合Mapper输出类型要求的局部最大值计算类。

额外建议

  • 运行作业时,通过yarn logs -applicationId <app-id>查看日志,能更清晰地定位错误。
  • 对输入数据进行合法性校验,避免因脏数据导致的异常。
  • 避免使用静态变量存储线程不安全的数据,尤其是在多线程执行的Map任务中。

内容的提问来源于stack exchange,提问作者ADNANE MEHDAOUI

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 21:12:12