跳到正文
Joeplover
学习笔记·2025-10-30·约 2 分钟阅读

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);
        }
    }
}

清洗后的增值

数据清洗不只是"清理脏数据",它会产生两层增值:

  1. 结构化:把文本描述转成可计算的字段(如"15k-25k"转为 min=15k, max=25k)
  2. 标准化:同一概念统一表达,方便聚合统计

实时 vs 离线清洗

方式适用场景工具
离线批量清洗历史数据全量清洗MapReduce、Spark
实时流清洗数据入库前实时过滤Flink、Kafka Streams
增量清洗每天新增数据的清洗Spark Streaming

在实际项目中,通常是离线 + 增量的结合——首次全量清洗历史数据,后续每天增量清洗新数据。