0%

数据输出

MapReduce OutputFormat 解析:结果怎么从键值对变成文件

MapReduce 处理完数据后,最后一步是把结果写出去。

OutputFormat 就是干这个的——把 Reduce 输出的 <key, value> 键值对,转成目标格式写到指定位置。

它干三件事:

  1. 格式转换:把键值对转成文本、二进制或自定义格式
  2. 写文件:写到 HDFS 指定路径
  3. 路径校验:检查输出目录是否存在、权限对不对

和 InputFormat 是对称的——输入从文件到键值对,输出从键值对到文件。

OutputFormat 的继承结构

1
2
3
4
5
6
OutputFormat(抽象类)
├─ FileOutputFormat(文件输出基类)
│ ├─ TextOutputFormat(默认,文本输出)
│ ├─ SequenceFileOutputFormat(二进制序列文件)
│ └─ 自定义 OutputFormat
└─ 非文件类 OutputFormat(如数据库、消息队列)

所有文件类 OutputFormat 都继承自 FileOutputFormat,它负责输出目录的创建和校验,子类只需实现”怎么写”的逻辑。

三种常用 OutputFormat

1. TextOutputFormat——默认,人类可读

TextOutputFormat 是默认输出格式,每条键值对输出一行:

1
key\tvalue

键和值之间用 Tab(\t)分隔。输出文件是 part-r-xxxxx(xxxxx 是 Reduce 任务编号)。

1
2
// 什么都不用配,默认就是这个
job.setOutputFormatClass(TextOutputFormat.class);

输出示例:

1
2
hello	100
world 200

适用场景: 绝大多数输出——需要人工查看的报表、统计结果、日志分析输出。没有特殊要求就用它。

2. SequenceFileOutputFormat——二进制,紧凑高效

SequenceFile 是 Hadoop 自带的二进制键值对格式,适合作为后续 MapReduce 作业的输入。

1
2
3
4
job.setOutputFormatClass(SequenceFileOutputFormat.class);
// 开启压缩
SequenceFileOutputFormat.setCompressOutput(job, true);
SequenceFileOutputFormat.setOutputCompressorClass(job, SnappyCodec.class);

输出特征:

  • 二进制文件,直接 cat 看不了
  • 支持压缩,存储效率高
  • 保留类型信息,读回来还是原来的 Writable 类型

适用场景: 多阶段 MapReduce 作业的中间结果。文本格式慢、占空间,中间层用 SequenceFile 更高效。

3. 自定义 OutputFormat——写 CSV、JSON 或任何格式

如果默认的不够用——比如你想输出 CSV、JSON,或者直接写数据库——就自己写。

实现步骤:

  1. 继承 FileOutputFormat,重写 getRecordWriter()
  2. 实现 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
// 自定义 RecordWriter
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 {
// 输出 CSV 格式:key,value
out.writeBytes(key.toString() + "," + value.get() + "\n");
}

@Override
public void close(TaskAttemptContext context) throws IOException {
out.close();
}
}

// 自定义 OutputFormat
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);
}
}

// Driver 里设置
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
// 1. 在 Driver 类中初始化 MultipleOutputs  
job.setOutputFormatClass(TextOutputFormat.class);
MultipleOutputs.addNamedOutput(job, "typeA", TextOutputFormat.class, Text.class, IntWritable.class);
MultipleOutputs.addNamedOutput(job, "typeB", TextOutputFormat.class, Text.class, IntWritable.class);

// 2. 在 Reducer 中使用 MultipleOutputs 输出
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 {
// 按 key 前缀拆分输出路径
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 (默认输出,如果有的话)

适用场景: 按条件拆分的报表、按类别归档的结果。

OutputFormat 怎么选

场景 用什么 原因
最终结果,需要人工查看 TextOutputFormat 默认,直接 cat 能看
中间结果,后续还要 MR 处理 SequenceFileOutputFormat 二进制高效、可压缩
按类别拆到不同目录 MultipleOutputs 一个作业输出多个路径
需要 CSV/JSON 或特殊格式 自定义 OutputFormat 没现成的就自己写

两个最常见的坑:

坑1:输出目录已存在,作业报错

MapReduce 不允许覆盖已有输出目录。要么手动删:

1
hdfs dfs -rm -r /output

要么代码里自动删:

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 数量,别盲目设很大。

欢迎关注我的其它发布渠道