0%

数据输入

MapReduce InputFormat 深度解析:数据怎么从文件变成键值对

MapReduce 处理数据的第一步,是把原始文件读进来、切成块、转成键值对,交给 Map 函数去处理。

这个过程的”总指挥”就是 InputFormat。它干三件事:

  1. 切分:把输入数据切成多个逻辑分片(InputSplit),每个分片交给一个 Map 任务
  2. 读取:每个分片通过 RecordReader 逐条读数据,转成 <key, value>
  3. 校验:检查输入路径对不对

InputSplit:逻辑分片,不是物理切块

InputSplit 是 MapReduce 的逻辑分片,记录了”读哪些数据”。

它不存储数据本身,只存三个信息:

  • 文件路径:读哪个文件
  • 起始位置:从文件的哪个字节开始读
  • 长度:读多少字节
  • 位置信息:数据在哪个 DataNode 上(用于任务本地化调度)
1
2
3
4
5
6
7
8
9
10
11
12
public abstract class InputSplit {  
// 获取分片大小(字节)
public abstract long getLength() throws IOException, InterruptedException;

// 获取分片所在的节点位置(DataNode 主机名)
public abstract String[] getLocations() throws IOException, InterruptedException;

// 获取更详细的位置信息(如机架信息),用于高级调度
public SplitLocationInfo[] getLocationInfo() throws IOException {
return null;
}
}

InputSplit vs HDFS Block:

InputSplit HDFS Block
是什么 逻辑分片,MapReduce 的概念 物理存储块,HDFS 的概念
存什么 不存数据,只存”位置+长度” 存实际数据
关系 一个 Split 通常对应一个 Block(默认) Block 是 Split 的数据来源

默认情况下,Split 大小 = Block 大小(128MB)。所以一个文件有几个 Block,就有几个 Map 任务。

FileInputFormat:分片大小的计算公式

FileInputFormat 是所有文件类 InputFormat 的父类,定义了分片的核心逻辑:

1
splitSize = max(minSize, min(maxSize, blockSize))

三个参数控制分片大小:

参数 默认值 作用
blockSize 128MB HDFS 块大小
minSize 1B 最小分片,设大了可以合并小文件
maxSize 无限大 最大分片,设小了可以拆分大文件

举个例子:

  • maxSize=64MB,分片=64MB → 一个 128MB Block 会被切成 2 个 Split,产生 2 个 Map 任务
  • minSize=256MB,分片=256MB → 两个 128MB Block 合并成 1 个 Split,产生 1 个 Map 任务

常用 InputFormat:不同场景用不同的

1. TextInputFormat——默认,按行读文本

最常用的输入格式,逐行读取文本文件,每行生成一个键值对:

  • key:该行在文件中的字节偏移量(LongWritable)
  • value:该行的文本内容(Text),不含换行符
1
2
输入:hello world\nhadoop mapreduce\n
输出:(0, "hello world")、(12, "hadoop mapreduce")

适用场景: 日志分析、词频统计等普通文本处理,没有特殊要求就用它。

2. KeyValueTextInputFormat——每行自带键值对

适合每行已经是 “key value” 格式的文件,按分隔符自动拆成键和值。

1
2
输入:name\tzhangsan\nage\t20
输出:("name", "zhangsan")、("age", "20")

默认分隔符是 Tab(\t),可以改成别的:

1
conf.set(KeyValueLineRecordReader.KEY_VALUE_SEPARATOR, ",");

适用场景: 配置文件、TSV/CSV 文件、已经按 key-value 组织的日志。

3. NLineInputFormat——按行数切分,不按大小

TextInputFormat 按字节大小切分,NLineInputFormat 按行数切分。每个 Map 任务处理固定行数。

1
文件 100 行,N=10 → 10 个 Map 任务,每个处理 10 
1
2
conf.setInt("mapreduce.input.lineinputformat.linespermap", 10);
job.setInputFormatClass(NLineInputFormat.class);

适用场景: 每行处理时间相近、希望均匀分配任务行数的场景。

4. CombineTextInputFormat——专门治小文件

小文件太多,每个文件一个 Map 任务,任务数量爆炸、调度开销大。CombineTextInputFormat 把多个小文件合并成一个 Split。

1
2
job.setInputFormatClass(CombineTextInputFormat.class);
CombineTextInputFormat.setMaxInputSplitSize(job, 20 * 1024 * 1024); // 20MB

多个小文件会按大小”拼”到一个 Split 里,直到接近 maxSplitSize

适用场景: 大量小文件(几 KB 到几十 MB),想减少 Map 任务数量。

5. 自定义 InputFormat——处理特殊格式

如果数据是二进制文件、图片、或者自定义协议格式,没有现成的 InputFormat,就自己写。

继承 FileInputFormat,重写 createRecordReader(),实现自己的 RecordReader 定义”怎么读、转成什么键值对”。

1
2
3
4
5
6
7
public class BinaryInputFormat extends FileInputFormat<LongWritable, BytesWritable> {
@Override
public RecordReader<LongWritable, BytesWritable> createRecordReader(
InputSplit split, TaskAttemptContext context) {
return new BinaryRecordReader();
}
}

四种 InputFormat 怎么选?

场景 用什么 原因
普通文本文件 TextInputFormat 默认,不用配置
已有 key-value 格式 KeyValueTextInputFormat 自动拆键值,省得自己 split
每行处理时间相近 NLineInputFormat 均匀分配行数
小文件太多 CombineTextInputFormat 合并分片,减少 Map 任务
二进制/自定义格式 自定义 InputFormat 没有现成的,自己实现

两个最常见的选型错误:

  • 小文件多还用 TextInputFormat → Map 任务爆炸,调度开销巨大 → 换 CombineTextInputFormat
  • 键值对文件用 TextInputFormat → Map 里自己 split 字符串,多写代码 → 换 KeyValueTextInputFormat

分片大小调优的几个场景

场景1:小文件多,想减少 Map 任务数

增大 minSize,让多个小文件合并成一个 Split:

1
2
3
4
<property>
<name>mapreduce.input.fileinputformat.split.minsize</name>
<value>134217728</value> <!-- 128MB -->
</property>

或者直接用 CombineTextInputFormat。

场景2:大文件单块太大,想多起几个 Map 并行处理

减小 maxSize,把一个 Block 切成多个 Split:

1
2
3
4
<property>
<name>mapreduce.input.fileinputformat.split.maxsize</name>
<value>67108864</value> <!-- 64MB -->
</property>

场景3:数据本地性差,Map 任务大量跨节点读数据

检查 InputSplit 的位置信息是否准确。如果分片跨 Block 边界,NameNode 返回的位置信息可能不准确,导致任务被调度到非本地节点。保持 splitSize = blockSize 能最大程度保证本地性。

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