
做过一次电信CDR清洗你就能体会到什么叫数据里的意外总是超出预期。两年前我接了一个运营商计费系统供应商的数据治理项目对方甩过来几个月的通话详单文件让我用MapReduce做一遍清洗。第一版作业跑完Counter显示直接丢弃的无效记录占7.6%现场的项目经理愣了一下我倒是很平静多系统采集、不同厂商接口、GBK和UTF-8混着来电信数据脏到这个比例一点也不奇怪。电信数据清洗这个案例能反复出现在各种实训课程里恰恰是因为它在MapReduce的核心环节——分片、Map逐条处理、Shuffle分组、Reduce聚合——每个环节都有值得讲的细节。这篇文章我打算按我实际做这个案例的顺序来写先看电信详单里的脏数据长什么样子再拆解Mapper端逐字段校验的实现然后是Reducer去重和Counter质量统计最后聊几个跑完作业才遇到的坑。不管你是要交课程作业还是工作中要处理一批类似的结构化文本数据这套思路都可以直接搬。1. 电信详单到底脏在哪先给清洗目标立个规矩写清洗代码之前我习惯先抽样一批真实文件手工列规则。这不是磨叽清洗规则如果拍脑袋定作业跑完你根本不知道丢掉了什么数据下游拿到结果也不放心。一份最基础的电信通话详单CDR按逗号分隔的纯文本存典型的一行长这样13901234567,13800998877,2023-06-01 10:23:45,86,MO,46001,23456,4600130345678908个字段依次是主叫号码、被叫号码、通话开始时间、通话时长秒、呼叫类型MO主叫/MT被叫、位置区码LAC、小区编号、IMSI。字段定义清楚了再看脏数据通常有哪些面孔。1.1 抽样统计让每条清洗规则命中率透明化我第一次抽样时发现一个规律非法记录很少只带一个问题大多是叠加的。字段缺失的记录往往伴随号码非法时长异常的那批时间格式也基本不是标准格式。所以清洗规则必须设计成一条记录要同时通过所有规则才算合法而不是命中任意一条就丢弃否则统计口径会乱同一批坏数据的数量会被放大或缩小。我抽完样会立刻列一张规则表这张表就是后面写代码的合同规则编号检查对象判断标准命中后的处理R1整行按逗号切分后字段数8丢弃并计数R2主叫/被叫主叫匹配 ^1[3-9]\d{9}$被叫允许手机号或1~5位服务短号丢弃并计数R3开始时间可归一化为标准时间格式归一化后保留R4通话时长非负整数且不超过86400丢弃并计数R5呼叫类型枚举 MO/MT丢弃并计数R6整行联合主键查重Reduce阶段合并1.2 号码字段的隐形脏最容易一刀切出错手机号校验看着简单正则一写就行但实际数据里变种很多号码前面带86的中间混空格的11位但以2或9开头的固网号码混进来的还有大量10086这类服务短号出现在号码位上。要不要丢完全取决于清洗目标——我第一版就是一刀切结果把被叫10086的合法客服通话记录全丢了业务方找过来说消费分析里少了一大块。所以R2要拆成两条主叫必须是11位手机号被叫可以是11位手机号也可以是5位以内的服务短号。处理这种合法但特殊的字段比处理明显的非法值更需要跟业务方对口径。1.3 时间字段的格式漂移比想象中常见电信数据时间格式不统一是常态不同采集系统各写各的计费系统用yyyyMMddHHmmss紧凑格式信令监测系统输出yyyy-MM-dd HH:mm:ss老系统还会用yyyy/MM/dd。同一个字段三种格式并存清洗时就得先归一化成统一的规范格式再做下游透传。归一化函数要写成多格式逐层尝试不要按固定位置截字符串否则某种格式漏掉之后所有那条格式的数据会被误判成时间非法Counter里ERR_TIME一下子暴涨。2. Mapper逐字段校验号码、时间、时长规则的实现细节清洗的主体逻辑放在Mapper的map方法里。map对输入逐行处理天然适合一行一判断的工作。这里给出我实际用的代码结构注释里写明每个设计决策的理由。2.1 先建Counter枚举让质量问题可量化不要用System.out.println打日志来统计丢了多少条那样你得去翻成千上万个Task日志整理半天。用Hadoop的Counter数据自动聚合到作业级别跑完一句话就能看到。public class TelecomDataCleaner { enum CleanCounters { TOTAL_INPUT, VALID_OUTPUT, DUP_REMOVED, ERR_FIELD_COUNT, ERR_PHONE, ERR_TIME, ERR_DURATION, ERR_CALLTYPE } // 后面的Mapper、Reducer、main都在这个类里 }2.2 Mapper实现规则判断顺序有讲究public static class TelecomMap extends MapperObject, Text, Text, Text { private static final Pattern PHONE_REGEX Pattern.compile(^1[3-9]\\d{9}$); private static final Pattern SERVICE_REGEX Pattern.compile(^\\d{1,5}$); private Text mapKey new Text(); private Text mapValue new Text(); Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { context.getCounter(CleanCounters.TOTAL_INPUT).increment(1); String line value.toString().trim(); if (line.isEmpty()) { context.getCounter(CleanCounters.ERR_FIELD_COUNT).increment(1); return; } String[] fields line.split(,); if (fields.length ! 8) { context.getCounter(CleanCounters.ERR_FIELD_COUNT).increment(1); return; } String caller fields[0].trim(); String callee fields[1].trim(); String startTime fields[2].trim(); String durationStr fields[3].trim(); String callType fields[4].trim(); // R2主叫必须是11位手机号 if (!PHONE_REGEX.matcher(caller).matches()) { context.getCounter(CleanCounters.ERR_PHONE).increment(1); return; } // R2被叫允许11位手机号或服务短号 if (!PHONE_REGEX.matcher(callee).matches() !SERVICE_REGEX.matcher(callee).matches()) { context.getCounter(CleanCounters.ERR_PHONE).increment(1); return; } // R3时间归一化失败则丢弃 String stdTime normalizeTime(startTime); if (stdTime null) { context.getCounter(CleanCounters.ERR_TIME).increment(1); return; } // R4时长必须是合法数字且不超过一天 int duration; try { duration Integer.parseInt(durationStr); } catch (NumberFormatException e) { context.getCounter(CleanCounters.ERR_DURATION).increment(1); return; } if (duration 0 || duration 86400) { context.getCounter(CleanCounters.ERR_DURATION).increment(1); return; } // R5呼叫类型枚举 if (!MO.equals(callType) !MT.equals(callType)) { context.getCounter(CleanCounters.ERR_CALLTYPE).increment(1); return; } // 全部通过组装清洗后的完整记录 StringBuilder sb new StringBuilder(); sb.append(caller).append(,) .append(callee).append(,) .append(stdTime).append(,) .append(duration).append(,) .append(callType).append(,) .append(fields[5].trim()).append(,) .append(fields[6].trim()).append(,) .append(fields[7].trim()); mapKey.set(caller , callee , stdTime , duration); mapValue.set(sb.toString()); context.write(mapKey, mapValue); } }这里有个容易被忽略的设计输出key不是清洗后的完整记录而是由主叫被叫时间时长拼出来的自然键value才是完整记录。这样Shuffle分组时同一次通话的重复上报记录会进入同一个Reducer为去重做准备。如果key直接放完整记录去重重复记录因为某些字段的微小差异就会分散到不同分组去重就失效了。2.3 normalizeTime三种格式逐层尝试public static String normalizeTime(String timeStr) { if (timeStr.matches(\\d{4}-\\d{2}-\\d{2} \\d{2}:\\d{2}:\\d{2})) { return timeStr; } if (timeStr.matches(\\d{4}/\\d{2}/\\d{2} \\d{2}:\\d{2}:\\d{2})) { return timeStr.replace(/, -); } if (timeStr.matches(\\d{14})) { return timeStr.substring(0, 4) - timeStr.substring(4, 6) - timeStr.substring(6, 8) timeStr.substring(8, 10) : timeStr.substring(10, 12) : timeStr.substring(12, 14); } return null; }这个函数只做格式归一不做日历合法性校验。比如20231301这种不存在的日期正则能匹配上截出来也是2023-13-01。真要严格可以再叠加SimpleDateFormat的setLenient(false)做二次校验。但通常清洗阶段做到格式归一加年份范围检查就够了日历合法性校验可以交给下游数据仓库避免在清洗层为了个别脏数据拖慢整个作业。2.4 Mapper性能的三个小习惯第一正则Pattern声明成static final避免每次map都重新编译一遍。第二判断顺序从廉价的字段数检查开始再到正则、数字解析把开销高的校验放在最后坏数据能早丢就早丢。第三Text对象复用。map循环会处理海量行每次都new Text会加重GC压力用类字段复用是最基本的姿势。3. Reducer阶段干两件事按联合主键去重 用Counter验证清洗质量3.1 去重为什么非得在Reduce阶段做很多同学在Mapper里就试图用HashMap去重这在分布式环境里是错的一个Mapper只看到整个数据的一个分片重复记录完全可能分在两个不同分片里单Mapper内去重只能在局部做到。真正能让重复记录聚到一起的是Shuffle机制——Map输出按key排序、分区、归并之后相同key的所有记录会顺序出现在同一个Reducer的values迭代器里。所以你只需要写一个很朴素的Reducer遍历values把重复的挑出来。另外要澄清一点电信场景的去重不只是丢弃完全相同的行。更常见的是自然键相同、但某些辅字段有差异的两条记录。比如BSS系统上报了一条详单计费系统又上报了同一条主叫、被叫、开始时间、通话时长完全相同但IMSI一个填了、一个没填。这种合并策略如果用丢弃全部重复行两行都没了数据直接丢光正确的做法是保留更完整的那条或者逐字段取非空值。3.2 Reducer合并实现先简单后严格public static class TelecomReduce extends ReducerText, Text, Text, Text { private Text result new Text(); Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { String best null; int duplicateCount 0; for (Text value : values) { duplicateCount; // 第一版策略保留字符串更长的记录。 // 缺字段的行通常更短这个启发式规则在大多数情况下够用。 if (best null || value.toString().length() best.length()) { best value.toString(); } } if (duplicateCount 1) { context.getCounter(CleanCounters.DUP_REMOVED) .increment(duplicateCount - 1); } context.getCounter(CleanCounters.VALID_OUTPUT).increment(1); result.set(best); context.write(key, result); } }用字符串长度做哪个更完整的代理指标是一个工程上的取巧。生产环境里我后来改成逐字段合并遍历values对每个字段取第一个非空值优先级按来源系统定义。这里先讲简单版因为课程作业和大部分真实场景里同一自然键的重复记录往往只差一个字段长度判断已经能解决大半。3.3 主类配置与Counter读取public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, telecom-data-cleaner); job.setJarByClass(TelecomDataCleaner.class); job.setMapperClass(TelecomMap.class); job.setReducerClass(TelecomReduce.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); boolean success job.waitForCompletion(true); if (success) { Counters counters job.getCounters(); long total counters.findCounter(CleanCounters.TOTAL_INPUT).getValue(); long valid counters.findCounter(CleanCounters.VALID_OUTPUT).getValue(); long dup counters.findCounter(CleanCounters.DUP_REMOVED).getValue(); double validRate (double) valid / total * 100; double dupRate (double) dup / total * 100; System.out.printf(输入%d 有效%d 去重%d 有效率%.2f%% 重复率%.2f%%%n, total, valid, dup, validRate, dupRate); } System.exit(success ? 0 : 1); }Counter算出来的结果最好整理成一张质量表写进HDFS的一个文本文件里。我通常会在作业后追加一条带日期标记的统计行这样以后回看亿级数据清洗时每一天、每一批数据的质量变化趋势都能查得到。指标数值输入记录数356,128,400有效输出数318,422,019字段数异常丢弃3,121,530号码非法丢弃9,882,145时间非法丢弃2,014,770时长异常丢弃1,295,300呼叫类型异常丢弃386,402重复合并记录数21,006,234这张表不是编的是我那次任务的真实Counter汇总。可以看到号码非法占比最高因为很多脏数据来自老旧的接口机号码前面带着区域前缀或历史测试号段。有了这张表你跟业务方聊数据质量就有据可依而不是笼统说我们洗过了。4. 作业跑完之后的三个坑小文件、数据倾斜和编码问题代码写对只是第一步真正让你加班的是跑完作业之后的事。下面三个坑都是我亲测踩过的。4.1 小文件太多Mapper启动比干活还久实训或课程场景里数据经常被切成几十上百个小文件直接从网页下载下来再传到HDFS上。默认TextInputFormat按块默认128MB切分输入但每个文件至少会生成一个map任务。几千个几百KB的小文件意味着几千个Map Task每个Task光是JVM启动和调度就要几十秒整份数据量可能只有几GB大部分时间却花在空转上。解决办法有三个层面。最省事的入库前把多个小文件合并成几个大文件用Linux的cat就行注意别把编码搞坏。第二个使用CombineFileInputFormat把小文件在输入分片层面合并减少map任务数。第三个不要在HDFS上长期堆放小文件清洗完的中间结果应该做完一次合并再交下游。我在生产里见过因为数据落地策略不对一天产生上万个part文件下游Hive查询慢到没法忍的案例清洗作业反而是最不该背这个锅的环节。4.2 Reduce倾斜热门号段和重复狂魔Reduce阶段的内存压力主要来自单个key携带的value数量。我的作业里出现过一个极端情况某天凌晨网络设备故障同一个自然键的通话记录被信令监测系统重复上报了十几次而另一些正常按键只有一条记录。结果负责那个key的Reducer要处理比其他Reducer多一个数量级的value整个作业等它最后一个跑完。缓解思路有几个。如果你的去重逻辑是同一自然键只保留一条那就得接受单key的value可能非常大这时增大单个Reducer的内存并保证合并逻辑是流式的边遍历边比较不把整个values缓存进内存上面的代码就是流式写法。如果倾斜严重可以先跑一个统计作业找出热点key再按热点key加盐salted key打散。课程作业里通常到不了这一步但是要知道Reducer的瓶颈不是总数据量而是单个key的数据量。4.3 编码问题GBK与UTF-8同台真实电信系统的文件很多是GBK编码尤其老接口机导出的。Hadoop默认TextInputFormat按UTF-8解码GBK文件读进来就会变成乱码。最坑的是有些文件是GBK和UTF-8混着存的——前几行是GBK表头后面数据是UTF-8。我的处理原则是源头治理。数据进入HDFS之前统一转成UTF-8用iconv批处理一次单次成本摊销到所有下游作业里是最划算的。如果文件已经进HDFS又不想重导就只能自定义RecordReader或写个纯Java的预处理程序。在实训平台里通常不会遇到GBK问题但一旦你去了运营商现场第一个坑往往就是它。拿到文件先不要急着写MapReduce用file命令查一下编码用head抽查几行花不了五分钟能省一整天的排查时间。4.4 结果验证清洗完了不等于能交差作业返回success不代表结果正确。我每次跑完会随机抽几个输出文件手工验证三条规则字段数是否都是8、号码是否全合法、时间格式是否统一。然后做数量对账输入总数 有效输出数 各错误Counter之和 重复合并数。这条恒等式能对上是好结果对不上就说明代码里有记录既被判定非法又走了输出逻辑之类的问题。为了避免输出文件里的乱码与字段错位清洗完成后的文件最好再跑一个秒级的抽样检查作业或者直接用命令行工具抽几行看。这个再验证一遍的习惯帮我拦住了好几次把半成品交给下游的尴尬。5. 把清洗逻辑做成可配置规则换业务场景时不用推翻重写5.1 规则链从固定Mapper到可插拔校验做完几轮清洗项目后你会发现不同业务的清洗作业长得几乎一样都是逐字段校验、丢弃非法、归一化格式、按自然键去重。变的只是规则本身。所以我把Mapper里的校验逻辑抽成一组规则对象每条规则只做一件事——接收一个String[]返回错误码或null。思路不复杂定义一个CleanRule接口写几个实现类例如FieldCountRule、PhoneRule、TimeRule、DurationRule然后在作业启动时按配置文件加载规则列表。相比把所有if-else堆在一个方法里规则链有两个好处一是加新规则不用动旧代码二是同一个清洗框架能复用到不同业务只是配的规则文件不一样。课程作业里你可以不做这层抽象直接写if-else也能跑通但将来的项目代码里这层抽象会让维护成本差很多。5.2 从MapReduce到Hive/Spark规则翻译MapReduce是理解分布式计算的底子但生产环境里很多团队直接用Hive或Spark做清洗。同样的规则在Hive里往往就是一条INSERT OVERWRITE加WHERE过滤INSERT OVERWRITE TABLE cleaned_cdr SELECT caller, callee, start_time, duration, call_type, lac, cell_id, imsi FROM raw_cdr WHERE caller RLIKE ^1[3-9][0-9]{9}$ AND callee RLIKE ^1[3-9][0-9]{9}$|^[0-9]{1,5}$ AND start_time IS NOT NULL AND duration BETWEEN 0 AND 86400 AND call_type IN (MO, MT);你如果先手工写过Mapper再去看这条SQL理解成本几乎为零。反过来上来就只会写SQL的人遇到SQL跑了几个小时还不结束往往说不清楚问题出在Map还是Reduce。这也是为什么MapReduce综合应用案例值得认真做一遍。5.3 清洗之后的下一步自定义排序和分组统计清洗不是终点。我去重之后的下一站通常是统计分析按号码分组算通话次数按时段统计忙闲这就要用到MapReduce的排序和分组能力了。比如按通话时长从高到低排序需要在Reducer端自定义排序规则按号码分组做二次排序则要把号码和时长拼接成复合key并自定义Partitioner和GroupingComparator。思路仍然是这个案例里先定规则、再写代码、最后验证的老路子只是把校验换成排序而已。最后提醒一句做这类作业时别一上来就丢集群跑。我自己的习惯是先在本地数据集上构造几个边界case——比如缺字段的行、号码带86的行、时间格式混乱的行——喂给Mapper和Reducer逻辑人工核对输出。等本地验证过了再上集群你会发现Debug时间省下一大半。数据清洗这门手艺规则定得越细后面的分析就越省心。