如何用Hadoop MapReduce按预设键将键值对输出至不同文件?
基于MapReduce实现按预设键值动态拆分输出到多文件
核心思路
要实现按任意数量预设文件中的键,将输入键值对拆分到对应输出文件,关键在于动态构建键到输出文件的映射关系,并利用Hadoop的MultipleOutputs组件实现多文件输出,避免固定Reducer数量的限制。整体流程为:
- 预加载所有预设键文件,建立「键-输出文件名」的映射;
- Map阶段匹配输入键与映射,输出目标文件名和原键值对;
- Reduce阶段按文件名分组,将对应内容写入指定文件。
代码实现
1. Driver类(作业配置)
负责初始化作业、加载预设文件到分布式缓存、配置多文件输出:
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.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs; import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; import org.apache.hadoop.util.GenericOptionsParser; import org.apache.hadoop.filecache.DistributedCache; public class SplitByPresetKeysDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs(); // 参数校验:输入X路径、输出根目录、至少一个预设键文件 if (otherArgs.length < 3) { System.err.println("Usage: SplitByPresetKeysDriver <X-input-path> <output-root> <preset-file1> [preset-file2...]"); System.exit(2); } Job job = Job.getInstance(conf, "SplitByPresetKeys"); job.setJarByClass(SplitByPresetKeysDriver.class); job.setMapperClass(SplitMapper.class); job.setReducerClass(SplitReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // 将所有预设键文件添加到分布式缓存,并记录对应输出文件名 for (int i = 2; i < otherArgs.length; i++) { Path presetPath = new Path(otherArgs[i]); DistributedCache.addCacheFile(presetPath.toUri(), job.getConfiguration()); String presetFileName = presetPath.getName(); conf.set("output.name." + presetFileName, presetFileName + "_new"); } // 配置输入输出路径 FileInputFormat.addInputPath(job, new Path(otherArgs[0])); FileOutputFormat.setOutputPath(job, new Path(otherArgs[1])); // 注册多文件输出 MultipleOutputs.addNamedOutput(job, "dynamicOutput", TextOutputFormat.class, Text.class, Text.class); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
2. Mapper类(键匹配与分组)
在初始化阶段读取分布式缓存中的预设文件,构建映射;Map阶段处理输入,输出目标文件名和原键值对:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.filecache.DistributedCache; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.BufferedReader; import java.io.FileReader; import java.io.IOException; import java.util.HashMap; import java.util.Map; public class SplitMapper extends Mapper<LongWritable, Text, Text, Text> { private final Map<String, String> keyToOutputFile = new HashMap<>(); @Override protected void setup(Context context) throws IOException, InterruptedException { Configuration conf = context.getConfiguration(); Path[] cacheFiles = DistributedCache.getLocalCacheFiles(conf); if (cacheFiles != null) { for (Path file : cacheFiles) { String presetFileName = file.getName(); String outputFileName = conf.get("output.name." + presetFileName); // 读取预设文件中的所有键,构建映射 BufferedReader reader = new BufferedReader(new FileReader(file.toString())); String line; while ((line = reader.readLine()) != null) { line = line.trim(); if (!line.isEmpty()) { keyToOutputFile.put(line, outputFileName); } } reader.close(); } } } @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 拆分输入行(示例格式:"Abby,4") String[] parts = value.toString().split(",", 2); if (parts.length != 2) { return; // 跳过格式错误的行 } String inputKey = parts[0].trim(); String targetFile = keyToOutputFile.get(inputKey); if (targetFile != null) { // 输出键为目标文件名,值为原输入行 context.write(new Text(targetFile), value); } // 未匹配的行可选择输出到默认文件,此处省略 } }
3. Reducer类(多文件输出)
利用MultipleOutputs将同一文件名分组的内容写入对应文件:
import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs; import java.io.IOException; public class SplitReducer extends Reducer<Text, Text, Text, NullWritable> { private MultipleOutputs<Text, NullWritable> multipleOutputs; @Override protected void setup(Context context) throws IOException, InterruptedException { multipleOutputs = new MultipleOutputs<>(context); } @Override protected void reduce(Text targetFile, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 将所有匹配的行写入对应文件 for (Text line : values) { multipleOutputs.write("dynamicOutput", line, NullWritable.get(), targetFile.toString()); } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { multipleOutputs.close(); } }
运行说明
- 打包作业:将上述三个类打包成JAR文件(如
split-by-preset-keys.jar); - 上传文件到HDFS:将输入文件X、所有预设键文件(A、B、C...)上传到HDFS;
- 提交作业:执行以下命令(替换路径为实际HDFS路径):
hadoop jar split-by-preset-keys.jar SplitByPresetKeysDriver /input/X /output /preset/A /preset/B /preset/C
- 查看输出:输出根目录
/output下会生成A_new-r-00000、B_new-r-00000等文件,对应匹配的键值对内容。
关键注意事项
- 分布式缓存:用于高效分发预设键文件到所有计算节点,避免重复读取HDFS;
- 动态映射:支持任意数量的预设文件,无需修改代码即可新增输出文件;
- 格式兼容:若输入键值对格式不同,只需修改Mapper中的
split逻辑即可适配; - 未匹配处理:可在Mapper中添加逻辑,将未匹配的键值对输出到默认文件(如
default_new)。
内容的提问来源于stack exchange,提问作者coder_98_mt
相关产品推荐
相关产品推荐

