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

如何用Hadoop MapReduce按预设键将键值对输出至不同文件?

基于MapReduce实现按预设键值动态拆分输出到多文件

核心思路

要实现按任意数量预设文件中的键,将输入键值对拆分到对应输出文件,关键在于动态构建键到输出文件的映射关系,并利用Hadoop的MultipleOutputs组件实现多文件输出,避免固定Reducer数量的限制。整体流程为:

  1. 预加载所有预设键文件,建立「键-输出文件名」的映射;
  2. Map阶段匹配输入键与映射,输出目标文件名和原键值对;
  3. 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();
    }
}

运行说明

  1. 打包作业:将上述三个类打包成JAR文件(如split-by-preset-keys.jar);
  2. 上传文件到HDFS:将输入文件X、所有预设键文件(A、B、C...)上传到HDFS;
  3. 提交作业:执行以下命令(替换路径为实际HDFS路径):
hadoop jar split-by-preset-keys.jar SplitByPresetKeysDriver /input/X /output /preset/A /preset/B /preset/C
  1. 查看输出:输出根目录/output下会生成A_new-r-00000、B_new-r-00000等文件,对应匹配的键值对内容。

关键注意事项

  • 分布式缓存:用于高效分发预设键文件到所有计算节点,避免重复读取HDFS;
  • 动态映射:支持任意数量的预设文件,无需修改代码即可新增输出文件;
  • 格式兼容:若输入键值对格式不同,只需修改Mapper中的split逻辑即可适配;
  • 未匹配处理:可在Mapper中添加逻辑,将未匹配的键值对输出到默认文件(如default_new)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 09:08:11