如何在Hadoop-2.9的MapReduce中复制任务实现单任务双Mapper执行?
实现MapReduce任务复制(单任务双Mapper执行)的方案
针对你在Hadoop 2.9集群上的需求——让每条记录对应的任务由两个独立的Mapper实例执行(即每个任务运行两次),这里提供两种可行的实现思路,具体选择取决于你的实际场景:
方法一:自定义InputFormat生成重复Split(推荐,实现真正的双Mapper任务)
这种方法通过扩展Hadoop的FileInputFormat,为每条原始数据对应的Split生成完全相同的副本,这样YARN会为每个副本分配一个独立的Map任务,从而实现同一份数据被两个Mapper处理。
步骤1:编写自定义DuplicateInputFormat类
import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.InputSplit; import org.apache.hadoop.mapreduce.JobContext; import org.apache.hadoop.mapreduce.RecordReader; import org.apache.hadoop.mapreduce.TaskAttemptContext; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.input.LineRecordReader; import java.io.IOException; import java.util.ArrayList; import java.util.List; public class DuplicateInputFormat extends FileInputFormat<LongWritable, Text> { @Override public List<InputSplit> getSplits(JobContext job) throws IOException { // 获取原始的Split列表 List<InputSplit> originalSplits = super.getSplits(job); List<InputSplit> duplicatedSplits = new ArrayList<>(originalSplits.size() * 2); // 为每个原始Split添加两个副本 for (InputSplit split : originalSplits) { duplicatedSplits.add(split); duplicatedSplits.add(split); } return duplicatedSplits; } @Override public RecordReader<LongWritable, Text> createRecordReader(InputSplit split, TaskAttemptContext context) throws IOException, InterruptedException { // 使用默认的LineRecordReader读取数据,无需修改 return new LineRecordReader(); } }
步骤2:在Job配置中使用自定义InputFormat
在你的MapReduce驱动类中,替换默认的InputFormat为刚才编写的DuplicateInputFormat:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class YourJobDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "DuplicateTaskJob"); // 设置自定义InputFormat job.setInputFormatClass(DuplicateInputFormat.class); // 其他常规配置:设置Mapper、Reducer类,输出键值类型等 job.setJarByClass(YourJobDriver.class); job.setMapperClass(YourOriginalMapper.class); job.setOutputKeyClass(YourOutputKey.class); job.setOutputValueClass(YourOutputValue.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
注意事项
- 该方法会将Map任务总数翻倍,集群资源消耗也会相应增加,需要确保你的5个Slave节点有足够的资源承载额外的Map任务。
- 如果你的原始InputFormat是将单条记录作为一个Split(比如自定义的单记录Split InputFormat),那么每个记录会对应两个独立的Map任务,完全符合你的需求。
- 若原始Split包含多条记录,那么整个Split内的所有记录都会被两个Map任务处理,适合需要批量复制任务的场景。
方法二:在Mapper内部重复执行逻辑(适合无需独立Mapper实例的场景)
如果你的需求只是让每条记录的计算逻辑执行两次,不需要严格的两个独立Mapper任务,可以直接在Mapper的map方法中调用两次处理逻辑:
import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; public class YourDuplicateMapper extends Mapper<LongWritable, Text, YourOutputKey, YourOutputValue> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 第一次执行原始逻辑 processRecord(key, value, context); // 第二次执行相同逻辑 processRecord(key, value, context); } // 封装原始的记录处理逻辑 private void processRecord(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 这里写你原来的Mapper处理代码 // ... context.write(outputKey, outputValue); } }
这种方法的优势是无需修改InputFormat,实现简单,但缺点是两次计算在同一个Mapper进程内完成,并非真正的两个独立任务,无法利用集群的分布式资源并行执行两次计算。
内容的提问来源于stack exchange,提问作者mng97
相关产品推荐
相关产品推荐

