MapReduce 处理完数据后,最后一步是把结果写出去。
OutputFormat 就是干这个的——把 Reduce 输出的 <key, value> 键值对,转成目标格式写到指定位置。
它干三件事:
- 格式转换:把键值对转成文本、二进制或自定义格式
- 写文件:写到 HDFS 指定路径
- 路径校验:检查输出目录是否存在、权限对不对
和 InputFormat 是对称的——输入从文件到键值对,输出从键值对到文件。
1 2 3 4 5 6
| OutputFormat(抽象类) ├─ FileOutputFormat(文件输出基类) │ ├─ TextOutputFormat(默认,文本输出) │ ├─ SequenceFileOutputFormat(二进制序列文件) │ └─ 自定义 OutputFormat └─ 非文件类 OutputFormat(如数据库、消息队列)
|
所有文件类 OutputFormat 都继承自 FileOutputFormat,它负责输出目录的创建和校验,子类只需实现”怎么写”的逻辑。
1. TextOutputFormat——默认,人类可读
TextOutputFormat 是默认输出格式,每条键值对输出一行:
键和值之间用 Tab(\t)分隔。输出文件是 part-r-xxxxx(xxxxx 是 Reduce 任务编号)。
1 2
| job.setOutputFormatClass(TextOutputFormat.class);
|
输出示例:
适用场景: 绝大多数输出——需要人工查看的报表、统计结果、日志分析输出。没有特殊要求就用它。
SequenceFile 是 Hadoop 自带的二进制键值对格式,适合作为后续 MapReduce 作业的输入。
1 2 3 4
| job.setOutputFormatClass(SequenceFileOutputFormat.class);
SequenceFileOutputFormat.setCompressOutput(job, true); SequenceFileOutputFormat.setOutputCompressorClass(job, SnappyCodec.class);
|
输出特征:
- 二进制文件,直接 cat 看不了
- 支持压缩,存储效率高
- 保留类型信息,读回来还是原来的 Writable 类型
适用场景: 多阶段 MapReduce 作业的中间结果。文本格式慢、占空间,中间层用 SequenceFile 更高效。
如果默认的不够用——比如你想输出 CSV、JSON,或者直接写数据库——就自己写。
实现步骤:
- 继承
FileOutputFormat,重写 getRecordWriter()
- 实现
RecordWriter,定义 write()(怎么把键值对转成目标格式)和 close()(资源清理)
一个 CSV 输出的例子(简化版):
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35
| public class CsvRecordWriter extends RecordWriter<Text, IntWritable> { private FSDataOutputStream out; public CsvRecordWriter(FSDataOutputStream out) { this.out = out; } @Override public void write(Text key, IntWritable value) throws IOException { out.writeBytes(key.toString() + "," + value.get() + "\n"); } @Override public void close(TaskAttemptContext context) throws IOException { out.close(); } }
public class CsvOutputFormat extends FileOutputFormat<Text, IntWritable> { @Override public RecordWriter<Text, IntWritable> getRecordWriter(TaskAttemptContext context) throws IOException { Path outputPath = FileOutputFormat.getOutputPath(context); Path file = new Path(outputPath, "part-" + context.getTaskAttemptID().getTaskID().getId()); FileSystem fs = FileSystem.get(context.getConfiguration()); FSDataOutputStream out = fs.create(file); return new CsvRecordWriter(out); } }
job.setOutputFormatClass(CsvOutputFormat.class);
|
适用场景: 需要特定格式(CSV、JSON)、或者要写数据库(非文件系统)——但写数据库不建议直接每个 Record 一条 SQL,通常用批量提交或专用工具更合适。
MultipleOutputs:一个作业输出到多个路径
默认情况下,一个 MapReduce 作业只能输出到一个目录。
如果想把结果按类别分到不同目录——比如按地区、按日志级别——用 MultipleOutputs。
用法:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29
| job.setOutputFormatClass(TextOutputFormat.class); MultipleOutputs.addNamedOutput(job, "typeA", TextOutputFormat.class, Text.class, IntWritable.class); MultipleOutputs.addNamedOutput(job, "typeB", TextOutputFormat.class, Text.class, IntWritable.class);
public class MyReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private MultipleOutputs<Text, IntWritable> multipleOutputs;
@Override protected void setup(Context context) { multipleOutputs = new MultipleOutputs<>(context); }
@Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { if (key.toString().startsWith("A")) { multipleOutputs.write("typeA", key, values.iterator().next(), "output/typeA/"); } else { multipleOutputs.write("typeB", key, values.iterator().next(), "output/typeB/"); } }
@Override protected void cleanup(Context context) throws IOException, InterruptedException { multipleOutputs.close(); } }
|
生成的目录结构:
1 2 3 4
| /output/ ├── typeA-r-00000 ├── typeB-r-00000 └── part-r-00000 (默认输出,如果有的话)
|
适用场景: 按条件拆分的报表、按类别归档的结果。
| 场景 |
用什么 |
原因 |
| 最终结果,需要人工查看 |
TextOutputFormat |
默认,直接 cat 能看 |
| 中间结果,后续还要 MR 处理 |
SequenceFileOutputFormat |
二进制高效、可压缩 |
| 按类别拆到不同目录 |
MultipleOutputs |
一个作业输出多个路径 |
| 需要 CSV/JSON 或特殊格式 |
自定义 OutputFormat |
没现成的就自己写 |
两个最常见的坑:
坑1:输出目录已存在,作业报错
MapReduce 不允许覆盖已有输出目录。要么手动删:
要么代码里自动删:
1 2 3 4 5
| Path outputPath = new Path(args[1]); FileSystem fs = FileSystem.get(conf); if (fs.exists(outputPath)) { fs.delete(outputPath, true); }
|
坑2:Reduce 太多,输出太多小文件
每个 Reduce 任务生成一个输出文件。Reduce 数 = 100,输出文件 = 100。
如果输出文件都很小,NameNode 内存压力大。控制 Reduce 数量,别盲目设很大。