MapReduce InputFormat 深度解析:数据怎么从文件变成键值对
MapReduce 处理数据的第一步,是把原始文件读进来、切成块、转成键值对,交给 Map 函数去处理。
这个过程的”总指挥”就是 InputFormat。它干三件事:
- 切分:把输入数据切成多个逻辑分片(InputSplit),每个分片交给一个 Map 任务
- 读取:每个分片通过 RecordReader 逐条读数据,转成
<key, value> - 校验:检查输入路径对不对
InputSplit:逻辑分片,不是物理切块
InputSplit 是 MapReduce 的逻辑分片,记录了”读哪些数据”。
它不存储数据本身,只存三个信息:
- 文件路径:读哪个文件
- 起始位置:从文件的哪个字节开始读
- 长度:读多少字节
- 位置信息:数据在哪个 DataNode 上(用于任务本地化调度)
1 | public abstract class InputSplit { |
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 | 输入:hello world\nhadoop mapreduce\n |
适用场景: 日志分析、词频统计等普通文本处理,没有特殊要求就用它。
2. KeyValueTextInputFormat——每行自带键值对
适合每行已经是 “key value” 格式的文件,按分隔符自动拆成键和值。
1 | 输入:name\tzhangsan\nage\t20 |
默认分隔符是 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 | conf.setInt("mapreduce.input.lineinputformat.linespermap", 10); |
适用场景: 每行处理时间相近、希望均匀分配任务行数的场景。
4. CombineTextInputFormat——专门治小文件
小文件太多,每个文件一个 Map 任务,任务数量爆炸、调度开销大。CombineTextInputFormat 把多个小文件合并成一个 Split。
1 | job.setInputFormatClass(CombineTextInputFormat.class); |
多个小文件会按大小”拼”到一个 Split 里,直到接近 maxSplitSize。
适用场景: 大量小文件(几 KB 到几十 MB),想减少 Map 任务数量。
5. 自定义 InputFormat——处理特殊格式
如果数据是二进制文件、图片、或者自定义协议格式,没有现成的 InputFormat,就自己写。
继承 FileInputFormat,重写 createRecordReader(),实现自己的 RecordReader 定义”怎么读、转成什么键值对”。
1 | public class BinaryInputFormat extends FileInputFormat<LongWritable, BytesWritable> { |
四种 InputFormat 怎么选?
| 场景 | 用什么 | 原因 |
|---|---|---|
| 普通文本文件 | TextInputFormat | 默认,不用配置 |
| 已有 key-value 格式 | KeyValueTextInputFormat | 自动拆键值,省得自己 split |
| 每行处理时间相近 | NLineInputFormat | 均匀分配行数 |
| 小文件太多 | CombineTextInputFormat | 合并分片,减少 Map 任务 |
| 二进制/自定义格式 | 自定义 InputFormat | 没有现成的,自己实现 |
两个最常见的选型错误:
- 小文件多还用 TextInputFormat → Map 任务爆炸,调度开销巨大 → 换 CombineTextInputFormat
- 键值对文件用 TextInputFormat → Map 里自己 split 字符串,多写代码 → 换 KeyValueTextInputFormat
分片大小调优的几个场景
场景1:小文件多,想减少 Map 任务数
增大 minSize,让多个小文件合并成一个 Split:
1 | <property> |
或者直接用 CombineTextInputFormat。
场景2:大文件单块太大,想多起几个 Map 并行处理
减小 maxSize,把一个 Block 切成多个 Split:
1 | <property> |
场景3:数据本地性差,Map 任务大量跨节点读数据
检查 InputSplit 的位置信息是否准确。如果分片跨 Block 边界,NameNode 返回的位置信息可能不准确,导致任务被调度到非本地节点。保持 splitSize = blockSize 能最大程度保证本地性。