MapReduce 数据清洗:脏数据进来,干净数据出去,连 Reduce 都省了
数据清洗是数据分析的第一步,也是最重要的一步——“垃圾进,垃圾出”,数据不干净,后面的分析全是错的。
清洗就是三件事:
- 过滤掉坏数据(格式错、字段缺、值不对)
- 标准化好数据(日期统一、大小写统一)
- 去掉重复的
在 MapReduce 里做数据清洗有个特别的优势:只需要 Mapper,不需要 Reducer。
因为清洗是逐条处理——每条数据独立判断、独立转换,不需要聚合。省掉 Reduce 阶段,就省掉了 Shuffle,数据直接从 Mapper 写到 HDFS,快得多。
为什么数据清洗不需要 Reduce?
MapReduce 的标准流程是 Map → Shuffle → Reduce。但清洗场景不需要分组、不需要聚合。
每条数据独立处理,过滤掉坏的,转换好的,然后直接输出。
flowchart TD
A[原始数据HDFS] --> B[InputFormat 分片]
B --> C[Mapper 读取并清洗数据]
C -->|过滤/转换| D[符合条件的数据输出]
D --> E[OutputFormat 写入 HDFS]
所以配置就是一句话:
1 | job.setNumReduceTasks(0); |
Reduce 数量设成 0,Shuffle 就不会发生,数据直接走 Mapper → OutputFormat → HDFS。
省掉了排序、省掉了网络传输、省掉了磁盘溢写——对于只需要”过滤+转换”的作业,这是最大的性能优化。
一个例子:清洗 Nginx 日志
原始日志:
1 | 192.168.1.1 - [2023-10-01 12:00:00] "GET /index.html" 200 1024 |
清洗目标:
- 只保留状态码 200 的日志(过滤掉 404、500 等错误日志)
- 检查格式是否正确(字段数够不够)
Mapper 代码:
1 | public class LogCleanMapper extends Mapper<LongWritable, Text, Text, Text> { |
Driver:
1 | job.setMapperClass(LogCleanMapper.class); |
进阶:把脏数据也记下来
有时候光扔掉不够——你需要知道扔掉了多少、为什么扔。
用 MultipleOutputs 把干净数据和脏数据分开放:
1 | public class LogCleanMapper extends Mapper<...> { |
Driver 里注册两个输出:
1 | MultipleOutputs.addNamedOutput(job, "valid", TextOutputFormat.class, Text.class, Text.class); |
这样跑完你会得到两个目录:valid/ 是干净数据,invalid/ 是脏数据。可以回头分析脏数据的原因。
数据转换:把格式统一
除了过滤,数据清洗还经常做格式标准化。
比如日期格式不统一,有的是 2023-10-01,有的是 01/10/2023。清洗的时候统一成一种。
1 | // 假设输入格式:user_id,action,date |
数据转换的常见类型:
| 类型 | 示例 |
|---|---|
| 日期标准化 | 01/10/2023 → 2023-10-01 |
| 大小写统一 | "Active" / "active" → "active" |
| 去除首尾空格 | " hello " → "hello" |
| 缺失值补全 | 空字段 → "unknown" |
| 格式校验 | 邮箱、手机号正则校验 |
性能优化
1. 省掉 Reduce 是最重要的优化
setNumReduceTasks(0) 省掉了 Shuffle,这是数据清洗能跑得快的核心原因。
2. 开启压缩
清洗后的数据如果很大,开启输出压缩:
1 | FileOutputFormat.setCompressOutput(job, true); |
3. 调整 Mapper 内存
清洗作业是纯 Mapper 作业,内存主要给 Map 用:
1 | <property> |
4. 使用 CombineTextInputFormat 合并小文件
如果输入是大量小文件,用 CombineTextInputFormat 减少 Map 任务数:
1 | job.setInputFormatClass(CombineTextInputFormat.class); |