如何在Hadoop 2.9.0 MapReduce中实现一个Mapper对应一个文本文件
实现一个Mapper对应一个输入文件的Hadoop MapReduce方案(Hadoop 2.9.0)
Hey there! 针对你遇到的「让Mapper数量严格等于小文件数量」的需求,本质是要绕过Hadoop默认的小文件合并逻辑——因为Hadoop的FileInputFormat会把多个小文件打包成一个InputSplit,导致多个文件共用一个Mapper。下面我会一步步给你讲怎么用Java代码实现,以及需要扩展的关键类和方法。
核心思路
要实现一个文件对应一个Mapper,关键是让每个InputSplit只包含一个完整的文件,这样每个Mapper就只会处理一个文件。我们需要自定义InputFormat并调整Split生成逻辑,同时控制文件不可被切分。
需要扩展的类与方法
1. 扩展FileInputFormat(核心)
这是实现需求的关键,我们需要重写两个方法:
isSplitable():返回false,告诉Hadoop单个文件不能被切分,确保整个文件作为一个独立SplitgetSplits():遍历所有输入文件,为每个文件单独创建一个FileSplit,保证Split数量和文件数量完全一致
2. 可选:在Mapper中获取当前处理的文件信息
如果业务逻辑需要知道当前Mapper处理的是哪个文件,可以通过Context获取InputSplit并解析文件路径/名称。
完整代码实现示例
自定义InputFormat类
import org.apache.hadoop.fs.Path; 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 SingleFileInputFormat extends FileInputFormat<String, String> { // 禁止文件被切分,确保每个文件作为完整Split @Override protected boolean isSplitable(JobContext context, Path filename) { return false; } // 为每个输入文件生成独立的FileSplit @Override public List<InputSplit> getSplits(JobContext context) throws IOException { List<InputSplit> splits = new ArrayList<>(); List<Path> inputPaths = getInputPaths(context); for (Path path : inputPaths) { // 每个文件对应一个Split,起始位置0,长度为文件总大小 splits.add(new FileSplit(path, 0, getPathSize(context.getConfiguration(), path), null)); } return splits; } // 使用默认的LineRecordReader读取每行文本 @Override public RecordReader<String, String> createRecordReader(InputSplit split, TaskAttemptContext context) throws IOException { return new LineRecordReader(); } }
Job驱动类配置
在你的Job启动代码中,替换默认的InputFormat为我们自定义的类:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class SingleFilePerMapperJobDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "SingleFilePerMapperJob"); // 设置自定义InputFormat job.setInputFormatClass(SingleFileInputFormat.class); // 替换为你自己的Mapper、Reducer类 job.setMapperClass(YourBusinessMapper.class); job.setReducerClass(YourBusinessReducer.class); // 根据你的业务需求设置输出键值类型 job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // 设置输入输出路径 SingleFileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
Mapper中获取当前文件信息(可选)
如果需要在Mapper中标记数据来源文件,可以在setup方法中获取当前Split的文件信息:
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.lib.input.FileSplit; import java.io.IOException; public class YourBusinessMapper extends Mapper<Object, Text, Text, Text> { private String currentFileName; @Override protected void setup(Context context) throws IOException { // 从Context中获取当前Split,强转为FileSplit FileSplit split = (FileSplit) context.getInputSplit(); // 获取文件名 currentFileName = split.getPath().getName(); } @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { // 示例:将文件名作为输出键的一部分 context.write(new Text(currentFileName + "_key"), value); } }
注意事项
- 适配你的业务场景:代码中使用的
LineRecordReader是读取文本行的实现,如果你的文件是其他格式,可以替换为对应的RecordReader - 性能考量:因为你处理的是10-100个小文件,禁用切分不会有性能问题;如果是大文件,这种方式会导致单个Mapper处理压力过大,不适用
- API版本:代码基于Hadoop 2.x的
org.apache.hadoop.mapreduce新API,不要混用旧的org.apache.hadoop.mapred包API
内容的提问来源于stack exchange,提问作者cdt
相关产品推荐
相关产品推荐

