超全模态模型 × Harness 升级,升级解锁 ArkClaw,最新支持 DeepSeek-V4 系列与 GLM-5.1
你的MapReduce作业输出为空,核心原因是Mapper输出的Value类型不兼容Hadoop的序列化机制,同时存在Job配置缺失、Combiner使用错误等问题,导致Reducer无法接收到Mapper传输的数据。以下是具体问题分析和修复方案:
Hadoop要求Map和Reduce之间传输的键值对类型必须实现Writable接口(或其子接口),而你使用的ArrayList<Object>并非Hadoop的序列化类型,Map阶段输出的数据无法被正确序列化传输到Reduce阶段,导致Reducer完全收不到数据,自然没有输出。
Writable
ArrayList<Object>
在Main类中,你只配置了最终作业的输出类型(IntWritable、IntWritable),但没有声明Mapper的输出类型(Text和自定义Value类型)。Hadoop无法自动推断非标准Writable类型的序列化方式,进一步加剧了数据传输失败。
IntWritable
Text
你将Reducer3同时设置为Combiner和Reducer,但Combiner的输出类型必须与Mapper的输出类型一致(作为Reducer的输入),而当前Reducer3的输出是IntWritable和IntWritable,与Mapper的输出类型完全不匹配,导致Combiner阶段的数据无法传递到Reducer,甚至直接丢失数据。
Reducer3
Arbre类中的champs是静态数组,在Map任务的多线程执行环境中,多个线程会同时修改这个静态变量,导致部分记录解析错误,而你的catch块没有任何日志输出,无法发现这些错误。
Arbre
champs
创建实现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); } }
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时初始化。
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阶段输出最终结果。
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>
内容的提问来源于stack exchange,提问作者ADNANE MEHDAOUI
超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起
模型再升级,30秒超长叙事, 模态参考扩容
模型自由,工具不限,最新支持 Deepseek-V4 系列、GLM-5.3 系列
超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列
大模型19元起,Al应用9.9元畅享,新人首购爆款尽享优惠