MapReduce 数据清洗:招聘数据处理
Hadoop MapReduce 处理非结构化招聘数据,标准化薪资格式、技能标签和地区字段。
MapReduce 数据清洗:招聘数据集处理
数据清洗通常占数据分析工作量的 60% 以上。在我之前的实习中处理过 70 万条数据清洗任务,让我对这一点体会深刻。
招聘数据的问题
原始招聘数据通常包含以下问题:
| 问题 | 举例 |
|---|---|
| 空值 | 薪资字段为空 |
| 格式不一致 | "10k-15k" vs "10K-15K" vs "10000-15000" |
| 异常值 | 月薪 999999 |
| 重复数据 | 同一职位在多个平台重复抓取 |
| 乱码 | 编码问题导致的特殊字符 |
MapReduce 清洗流程
public class DataCleaning {
public static class CleanMapper
extends Mapper<Object, Text, Text, NullWritable> {
private Text outKey = new Text();
private int totalLines = 0;
private int filteredLines = 0;
public void map(Object key, Text value, Context context)
throws IOException, InterruptedException {
totalLines++;
String line = value.toString().trim();
// 1. 过滤空行
if (line.isEmpty()) return;
String[] fields = line.split(",");
// 2. 字段数量校验
if (fields.length != EXPECTED_FIELDS) return;
// 3. 统一薪资格式:"10k-15k" → 统一转小写
fields[3] = normalizeSalary(fields[3]);
// 4. 去除前后空格
for (int i = 0; i < fields.length; i++) fields[i] = fields[i].trim();
// 5. 过滤异常薪资
try {
double minSalary = parseMinSalary(fields[3]);
if (minSalary < 2000 || minSalary > 100000) return;
} catch (NumberFormatException e) {
return; // 字段值无法解析
}
// 6. 去重(以职位名+公司名为唯一标识)
outKey.set(fields[1] + "|" + fields[2]);
context.write(outKey, NullWritable.get());
filteredLines++;
}
protected void cleanup(Context context) {
System.out.println("清洗率: " + (totalLines - filteredLines) + "/" + totalLines);
}
}
}
清洗后的增值
数据清洗不只是"清理脏数据",它会产生两层增值:
- 结构化:把文本描述转成可计算的字段(如"15k-25k"转为 min=15k, max=25k)
- 标准化:同一概念统一表达,方便聚合统计
实时 vs 离线清洗
| 方式 | 适用场景 | 工具 |
|---|---|---|
| 离线批量清洗 | 历史数据全量清洗 | MapReduce、Spark |
| 实时流清洗 | 数据入库前实时过滤 | Flink、Kafka Streams |
| 增量清洗 | 每天新增数据的清洗 | Spark Streaming |
在实际项目中,通常是离线 + 增量的结合——首次全量清洗历史数据,后续每天增量清洗新数据。