电影票房数据清洗这件事,听起来像是大学里的大作业,但真正在Hadoop生态里跑过一遍之后,你会发现“脏数据”这三个字的分量。最近我在做一组电影票房历史数据的离线分析,数据来源是某票务平台导出的CSV,大概三十多万行。用pandas拉下来本地看了眼,结果差点劝退自己:日期有三种格式,票房字段有的带“亿”、有的带“万”、有的干脆是纯数字,片名后面跟着一堆HTML标签,还有同一条电影ID对应两个不同片名的记录。这种数据扔给任何统计脚本都是灾难,于是我把整条清洗链路搬到了MapReduce上,做了一个专门处理票房数据的清洗任务。
这篇文章我会把这套任务从环境准备、清洗规则设计、MapReduce代码实现到运行调优完完整整拆开讲,包括我踩过的坑和最终沉淀下的经验。如果你正准备做类似的数据预处理、MapReduce编程实训,或者手里也有一批不大不小的结构化数据要清洗,这篇内容可以直接拿去参考。
1. 为什么要给票房数据做清洗
1.1 一份脏数据到底能有多脏
做数据清洗之前,先要搞清楚原料有多“脏”。票务平台导出的数据表面上是规整的CSV,实际上每一列都可能有意外。我这次拿到的原始数据是这样的:
| movie_id | movie_name | type | country | release_date | box_office | rating | rating_people | director | actor |
|---|---|---|---|---|---|---|---|---|---|
| MOV_001 | 流浪地球2 | 科幻/冒险 | 中国内地 | 2023-01-22 | 40.29亿 | 8.3 | 1,023,430 | 郭帆 | 吴京 |
| MOV_002 | 满江红 | 喜剧/悬疑 | 中国大陆 | 2023/1/22 | 454400万 | 7.0 | 856120 | 张艺谋 | 沈腾/易烊千玺 |
| MOV_003 | 我不是药神 | 剧情 | 中国 | 2018年7月5日 | 30.7 | 9.0 | 1617898 | 文牧野 | 徐峥 |
| MOV_001 | 流浪地球2 | 科幻 冒险 | 中国内地 | 20230122 | 402900.0 | 8.3分 | 1,023,430 | 郭帆 | 吴京 |
| MOV_004 | 满江红(重映) | 喜剧悬疑 | China | 2023-01-22 | 45.44 | 7.0 | 856120 | 张艺谋 | 沈腾,易烊千玺 |
这里面有几类特别常见的问题:日期格式三套并行,同一天上映能写出三种写法;票房字段有的带“亿”字、有的带“万”字、有的直接是万元数值,单位不统一;评分后面跟着“分”这种中文说明;评分人数里带了千分位逗号;类型字段有的用斜线、有的用空格、有的用顿号;国家字段中英文混用;还有同一条电影记录的ID相同、片名却不同的重复数据。这些问题单独看都不算致命,但一旦进入聚合统计分析,比如按年份算总票房、按类型算平均评分,就会直接得出错误结论。
1.2 清洗的目标:从“能看”到“能用”
清洗不是简单地把空值填上,而是要让数据达到“可直接入数仓、可参与聚合计算”的状态。我给自己定的清洗目标很明确:字段能对齐、格式能统一、记录不重复、异常剔除有依据、输出结构标准化。
具体拆成五条:第一,编码统一,所有乱码、非法字符清理干净;第二,字段对齐,错位的记录要么修复要么剔除;第三,日期、票房、评分数值统一成同一种格式和单位;第四,按电影维度去重,同名、同ID的记录合并;第五,输出结果可以直接加载到Hive或后续分析脚本里。这五条是清洗任务的验收标准,缺一条后面分析阶段就要还债。
还有一个容易被忽略的目标:清洗过程中要能输出质量报告。也就是要知道这一轮清洗到底过滤了多少行、修复了多少字段、哪类脏数据占比最高。没有这些数字,清洗就是一笔糊涂账。后面我会说怎么用MapReduce的Counters来做这件事。
1.3 为什么选 MapReduce 而不是 Pandas
聊到数据清洗,90%的人第一反应是用pandas。我也一样,本地写脚本确实方便,但是这次的情况不一样。第一,数据量虽然只有三十多万行,但分析链路是放在Hadoop集群上的,上游数据进了HDFS,下游还有Hive数据仓库,中间用一套本地Python脚本反而得来回倒数据。第二,MapReduce适合这种“逐行处理、按Key聚合”的场景,清洗的本质就是把每一行读进来、按规则修正、再按某个维度合并去重,这和MapReduce的编程模型天然匹配。第三,如果未来数据量从三十万涨到三千万,pandas单机处理就会吃力,而同一份MapReduce代码几乎不用改就能扛住更大的数据规模。
当然,MapReduce也有它的笨重之处,比如调试周期比Python长、代码量也更大。所以我在设计时把“可维护性”放在了重要位置,用自定义Writable封装电影记录,把每条清洗规则做成独立方法。这样清洗逻辑清晰,后面想加规则也方便。
2. 数据建模与清洗规则设计
2.1 字段口径与目标结构
清洗之前先定义目标表的字段结构。我把原始数据标准化成11个字段,每个字段都有明确的类型和口径:
| 字段名 | 类型 | 说明 |
|---|---|---|
| movie_id | String | 电影唯一标识,去重基准之一 |
| movie_name | String | 电影名称,去除广告后缀、HTML标签 |
| type | String | 类型,多个类型用“/”拼接 |
| country | String | 制片国家/地区,统一中文表示 |
| release_date | String | 上映日期,统一 yyyy-MM-dd |
| box_office | Double | 票房,统一单位为万元 |
| rating | Double | 评分,去“分”字,统一保留1位小数 |
| rating_people | Long | 评分人数,去千分位逗号 |
| director | String | 导演,非空判断 |
| actor | String | 主演,多个主演统一用逗号分隔 |
| source_status | String | 清洗标记,正常/修复/待人工确认 |
source_status 这个字段是我自己加的,它用来标记一条记录是原始就没问题,还是经过了修复,还是存在无法确定的异常需要人工确认。这样清洗结果不只是“对或错”,还能追溯每一条数据的处理过程。这个设计在数据分析阶段特别好用,看到某数据异常时可以直接回看来源。
2.2 清洗规则清单
规则不在多,而在能被明确执行。我整理了一张清洗规则表,直接对应代码里的每一条判断逻辑:
| 问题类型 | 示例 | 处理规则 |
|---|---|---|
| 空行、注释行 | 无内容或以#开头 | 直接过滤,不进入后续处理 |
| 分隔符异常 | Tab、全角逗号、多空格 | 先统一为半角逗号,再按CSV引号规则切分 |
| 字段缺失 | 类型为空 | 字段为空置为“未知”,关键字段(movie_id、movie_name)为空直接剔除 |
| 字段错位 | 年份写进类型栏 | 按字段类型校验,修复失败则标记为待人工确认 |
| 日期格式不统一 | 2023-01-22 / 2023/1/22 / 2018年7月5日 / 20230122 | 统一转为 yyyy-MM-dd |
| 数值含中文单位 | 40.29亿、454400万 | 亿统一乘10000换算为万元 |
| 数值含附加字符 | 8.3分、1,023,430 | 去“分”字、去千分位逗号 |
| 类型分隔符混乱 | 科幻 冒险 / 喜剧悬疑 | 统一为“/”分隔 |
| 重复记录 | 同ID同名多条 | 按 movie_id + movie_name 归并,合并字段 |
| 异常值 | box_office <= 0 | 保留但标记,暂不剔除 |
| 乱码与HTML标签 | 片名带 标签 | 正则匹配剔除 |
这套规则的核心思想是“能修则修,不能修则不盲目丢弃”。清洗的目的是让数据变得更干净,而不是把可疑数据一刀切掉。比如票房为0或负数的记录,有可能是数据缺失,也有可能是特殊发行方式导致的空档期记录,直接删了太可惜,先保留并打上标记,后面分析时灵活处理。
2.3 关键决策:Key和Value如何设计
MapReduce的清洗任务中,最关键的建模决策是map端输出什么Key、什么Value。这个决策直接决定了去重逻辑的复杂度。
我采用的方案是:Key用movie_id + "|" + movie_name,Value用自定义Writable对象MovieWritable封装整条记录的11个字段。
为什么Key要同时包含movie_id和movie_name?纯用movie_id会出问题:有些电影ID在数据源里被错误复用了,导致两部不同的电影被硬合并成一条。纯用movie_name也不行,重名电影太常见了。所以两个都带上,把“同ID且同片名”作为去重的基准。至于同一ID不同片名的情况,比如“满江红”和“满江红(重映)”,我单独设计了一套合并优先级,后面在Reducer部分详细说。
Value用自定义Writable而不是直接用Text好处很明显:第一,字段是强类型的,解析和校验逻辑可以封装在类里;第二,Reducer中要做字段合并,直接用对象比反复解析字符串省事得多;第三,代码可读性强,看类名就知道这是电影票房记录。
3. MapReduce 清洗任务实现
3.1 自定义Writable与CSV解析
先说MovieWritable。它实现了WritableComparable接口,封装11个字段。write和readFields的顺序必须完全一致,这是Hadoop序列化的硬要求,顺序不一致会导致反序列化后字段错乱,这种问题排查起来特别隐蔽。
import org.apache.hadoop.io.WritableComparable; import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; public class MovieWritable implements WritableComparable<MovieWritable> { private String movieId; private String movieName; private String type; private String country; private String releaseDate; private double boxOffice; private double rating; private long ratingPeople; private String director; private String actor; private String sourceStatus; public MovieWritable() { this.movieId = ""; this.movieName = ""; this.type = ""; this.country = ""; this.releaseDate = ""; this.boxOffice = 0.0; this.rating = 0.0; this.ratingPeople = 0L; this.director = ""; this.actor = ""; this.sourceStatus = "normal"; } @Override public void write(DataOutput out) throws IOException { out.writeUTF(movieId); out.writeUTF(movieName); out.writeUTF(type); out.writeUTF(country); out.writeUTF(releaseDate); out.writeDouble(boxOffice); out.writeDouble(rating); out.writeLong(ratingPeople); out.writeUTF(director); out.writeUTF(actor); out.writeUTF(sourceStatus); } @Override public void readFields(DataInput in) throws IOException { this.movieId = in.readUTF(); this.movieName = in.readUTF(); this.type = in.readUTF(); this.country = in.readUTF(); this.releaseDate = in.readUTF(); this.boxOffice = in.readDouble(); this.rating = in.readDouble(); this.ratingPeople = in.readLong(); this.director = in.readUTF(); this.actor = in.readUTF(); this.sourceStatus = in.readUTF(); } @Override public int compareTo(MovieWritable o) { int cmp = this.movieId.compareTo(o.movieId); if (cmp != 0) return cmp; return this.movieName.compareTo(o.movieName); } // 省略 getter/setter }CSV解析是一个细节很多的地方。直接按逗号split有一个致命问题:如果某个字段本身包含逗号,比如演员表写成“吴京,刘德华”,但这一列在原始CSV里用了引号包裹,那split就会把字段拦腰切断。我手写了一个处理引号的CSV解析方法:
import java.util.ArrayList; import java.util.List; public class CsvParseUtil { public static List<String> parseLine(String line) { List<String> fields = new ArrayList<>(); StringBuilder sb = new StringBuilder(); boolean inQuotes = false; for (int i = 0; i < line.length(); i++) { char c = line.charAt(i); if (inQuotes) { if (c == '"') { if (i + 1 < line.length() && line.charAt(i + 1) == '"') { sb.append('"'); i++; } else { inQuotes = false; } } else { sb.append(c); } } else { if (c == '"') { inQuotes = true; } else if (c == ',') { fields.add(sb.toString().trim()); sb.setLength(0); } else { sb.append(c); } } } fields.add(sb.toString().trim()); return fields; } }这个解析器虽然精简,但能处理绝大多数CSV引号和逗号场景。读这一段的重点是理解:数据清洗里最简单的“按逗号拆分”都有这么多讲究,真实数据远比教科书示例凶险。
3.2 CleanMapper:逐条清洗逻辑
Mapper负责把每行原始记录读进来,做字段级清洗,然后输出Text -> MovieWritable。清洗方法拆成独立函数,每个函数只处理一种脏数据,这样代码好维护,也方便加单测。
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; import java.util.List; import java.util.regex.Pattern; public class CleanMapper extends Mapper<Object, Text, Text, MovieWritable> { private Text outKey = new Text(); private MovieWritable outValue = new MovieWritable(); private static final Pattern HTML_TAG_PATTERN = Pattern.compile("<[^>]+>"); private static final Pattern ILLEGAL_CHAR_PATTERN = Pattern.compile("[\\u0000-\\u001f\u007f]"); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); if (line == null || line.trim().isEmpty()) { return; } String trimmed = line.trim(); if (trimmed.startsWith("#")) { return; } context.getCounter("CleanStat", "RAW_LINES").increment(1); List<String> fields; try { fields = CsvParseUtil.parseLine(trimmed); } catch (Exception e) { context.getCounter("CleanStat", "PARSE_ERROR").increment(1); return; } if (fields.size() < 10) { context.getCounter("CleanStat", "LESS_FIELDS").increment(1); return; } String movieId = cleanId(fields.get(0)); String movieName = cleanName(fields.get(1)); if (movieId.isEmpty() || movieName.isEmpty()) { context.getCounter("CleanStat", "NULL_KEY_FIELDS").increment(1); return; } outValue.setMovieId(movieId); outValue.setMovieName(movieName); outValue.setType(cleanJoinField(fields.get(2), "未知")); outValue.setCountry(cleanCountry(fields.get(3))); outValue.setReleaseDate(cleanDate(fields.get(4))); outValue.setBoxOffice(cleanBoxOffice(fields.get(5))); outValue.setRating(cleanRating(fields.get(6))); outValue.setRatingPeople(cleanRatingPeople(fields.get(7))); outValue.setDirector(fields.get(8).trim()); outValue.setActor(cleanActor(fields.get(9))); outKey.set(movieId + "|" + movieName); context.write(outKey, outValue); } private String cleanId(String raw) { if (raw == null) return ""; String cleaned = raw.trim(); if (cleaned.equalsIgnoreCase("null") || cleaned.equals("-")) { return ""; } return cleaned; } private String cleanName(String raw) { if (raw == null) return ""; String cleaned = HTML_TAG_PATTERN.matcher(raw).replaceAll(""); cleaned = ILLEGAL_CHAR_PATTERN.matcher(cleaned).replaceAll(""); cleaned = cleaned .replace("(重映)", "") .replace("(重映)", "") .trim(); return cleaned; } private String cleanJoinField(String raw, String defaultValue) { if (raw == null || raw.trim().isEmpty()) return defaultValue; String cleaned = raw.trim() .replace(",", "/") .replace(",", "/") .replace("、", "/") .replace(" ", "/") .replace("/", "/"); while (cleaned.contains("//")) { cleaned = cleaned.replace("//", "/"); } return cleaned; } private String cleanCountry(String raw) { if (raw == null || raw.trim().isEmpty()) return "未知"; String cleaned = raw.trim(); switch (cleaned) { case "China": case "CN": case "中国大陆": case "中国内地": return "中国内地"; case "中国香港": case "Hong Kong": return "中国香港"; case "中国台湾": case "Taiwan": case "中国台湾省": return "中国台湾"; default: return cleaned; } } private String cleanDate(String raw) { if (raw == null) return ""; String cleaned = raw.trim(); if (cleaned.isEmpty()) return "1970-01-01"; cleaned = cleaned.replace("年", "-").replace("月", "-").replace("日", ""); cleaned = cleaned.replace("/", "-").replace(".", "-"); cleaned = cleaned.replaceAll("\\s+", ""); String[] parts = cleaned.split("-"); if (parts.length != 3) { context.getCounter("CleanStat", "DATE_FORMAT_ERROR").increment(1); return "1970-01-01"; } String year = String.format("%04d", Integer.parseInt(parts[0])); String month = String.format("%02d", Integer.parseInt(parts[1])); String day = String.format("%02d", Integer.parseInt(parts[2])); return year + "-" + month + "-" + day; } private double cleanBoxOffice(String raw) { if (raw == null || raw.trim().isEmpty()) return 0.0; String cleaned = raw.trim(); double result = 0.0; try { if (cleaned.endsWith("亿")) { result = Double.parseDouble(cleaned.replace("亿", "")) * 10000; } else if (cleaned.endsWith("万")) { result = Double.parseDouble(cleaned.replace("万", "")); } else if (cleaned.endsWith("元")) { result = Double.parseDouble(cleaned.replace("元", "")) / 10000; } else { result = Double.parseDouble(cleaned); } } catch (NumberFormatException e) { context.getCounter("CleanStat", "BOX_OFFICE_PARSE_ERROR").increment(1); return 0.0; } if (result < 0) { context.getCounter("CleanStat", "NEGATIVE_BOX_OFFICE").increment(1); } return result; } private double cleanRating(String raw) { if (raw == null || raw.trim().isEmpty()) return 0.0; String cleaned = raw.trim().replace("分", ""); try { return Double.parseDouble(cleaned); } catch (NumberFormatException e) { context.getCounter("CleanStat", "RATING_PARSE_ERROR").increment(1); return 0.0; } } private long cleanRatingPeople(String raw) { if (raw == null || raw.trim().isEmpty()) return 0L; String cleaned = raw.trim().replace(",", "").replace(",", ""); try { return Long.parseLong(cleaned); } catch (NumberFormatException e) { context.getCounter("CleanStat", "RATING_PEOPLE_PARSE_ERROR").increment(1); return 0L; } } private String cleanActor(String raw) { if (raw == null || raw.trim().isEmpty()) return "未知"; String cleaned = raw.trim() .replace("/", ",") .replace("、", ",") .replace(",", ","); String[] actors = cleaned.split(","); StringBuilder sb = new StringBuilder(); for (String actor : actors) { String a = actor.trim(); if (!a.isEmpty()) { if (sb.length() > 0) sb.append(","); sb.append(a); } } return sb.toString(); } }这里有几个经验点值得展开说。
第一,日期解析不要迷信某一个格式。我在cleanDate里先统一把“年/月/日”替换成“-”,再把各种分隔符统一替换,最后用split拆三段后格式化补零。这样“2018年7月5日”和“2018/7/5”都能正确转换为“2018-07-05”。如果遇到固定格式以外的数据,也不要抛异常,记录计数后置为默认值,保证任务不因一条脏数据整体失败。
第二,票房单位的换算必须非常谨慎。我的统一口径是“万元”,因为票务平台导出时大部分数值都以万为单位了,少数大热影片用“亿”表示。40.29亿等于402900万,这个换算如果写错,后面所有汇总数据都会偏差巨大。我在这儿加了一个专门的计数器,清理每一条“亿”级数据时都可以追踪。
第三,MapReduce里的Counter是个好东西,但别滥用。我统计了RAW_LINES、PARSE_ERROR、LESS_FIELDS、NULL_KEY_FIELDS、DATE_FORMAT_ERROR等几类关键计数。这些计数在任务结束时会统一输出,相当于自动生成了一份数据质量报告。
3.3 CleanReducer:去重与字段合并
Reducer负责按Key归并同一条电影的相关记录,然后做字段级合并。合并规则我定义为:对于字符串字段,取非空且更长的那个;对于数值字段,取非空且看起来更合理的值;source_status字段用来标记这条记录是否经过修复。
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class CleanReducer extends Reducer<Text, MovieWritable, Text, Text> { @Override protected void reduce(Text key, Iterable<MovieWritable> values, Context context) throws IOException, InterruptedException { MovieWritable merged = new MovieWritable(); boolean hasValue = false; boolean repaired = false; for (MovieWritable val : values) { hasValue = true; merged.setMovieId(mergeString(merged.getMovieId(), val.getMovieId(), true)); merged.setMovieName(mergeString(merged.getMovieName(), val.getMovieName(), true)); String type = mergeString(merged.getType(), val.getType(), false); if (!type.equals(merged.getType())) repaired = true; merged.setType(type); merged.setCountry(mergeString(merged.getCountry(), val.getCountry(), false)); merged.setReleaseDate(mergeDate(merged.getReleaseDate(), val.getReleaseDate())); double box = mergeBoxOffice(merged.getBoxOffice(), val.getBoxOffice()); if (box != merged.getBoxOffice()) repaired = true; merged.setBoxOffice(box); double rating = mergeRating(merged.getRating(), val.getRating()); if (rating != merged.getRating()) repaired = true; merged.setRating(rating); long ratingPeople = mergeRatingPeople(merged.getRatingPeople(), val.getRatingPeople()); if (ratingPeople != merged.getRatingPeople()) repaired = true; merged.setRatingPeople(ratingPeople); merged.setDirector(mergeString(merged.getDirector(), val.getDirector(), false)); merged.setActor(mergeString(merged.getActor(), val.getActor(), false)); } if (!hasValue) { return; } merged.setSourceStatus(repaired ? "repaired" : "normal"); context.getCounter("CleanResult", "OUTPUT_RECORDS").increment(1); if (repaired) { context.getCounter("CleanResult", "REPAIRED_RECORDS").increment(1); } context.write(key, new Text(merged.toString())); } private String mergeString(String oldVal, String newVal, boolean keepLonger) { if (oldVal == null || oldVal.isEmpty()) return newVal; if (newVal == null || newVal.isEmpty()) return oldVal; if (keepLonger) { return oldVal.length() >= newVal.length() ? oldVal : newVal; } return oldVal; } private String mergeDate(String oldVal, String newVal) { if (oldVal == null || oldVal.isEmpty() || oldVal.equals("1970-01-01")) return newVal; if (newVal == null || newVal.isEmpty() || newVal.equals("1970-01-01")) return oldVal; return oldVal; } private double mergeBoxOffice(double oldVal, double newVal) { if (oldVal == 0.0) return newVal; if (newVal == 0.0) return oldVal; return Math.max(oldVal, newVal); } private double mergeRating(double oldVal, double newVal) { if (oldVal == 0.0) return newVal; if (newVal == 0.0) return oldVal; return (oldVal + newVal) / 2.0; } private long mergeRatingPeople(long oldVal, long newVal) { if (oldVal == 0L) return newVal; if (newVal == 0L) return oldVal; return Math.max(oldVal, newVal); } }合并逻辑中最难的是“同一ID但不同名称”的记录。我采取的方案是:在Mapper端已经通过cleanName把“(重映)”等后缀去掉了,所以“满江红(重映)”会变成“满江红”,和正常记录归到同一个Key下。对于完全无法对齐的名称,比如“老炮儿”和“老炮”,只能靠人工规则补。我在实际项目中维护了一张别名映射表,后续版本可以做成分布式缓存文件让Mapper加载,这次先不做过度设计。
另外注意,我并没有在Reducer里使用Combiner。原因是合并逻辑不是简单的加法或取最大值,它涉及到“非空优先”“保留较长字符串”“两个评分取平均”这类语义,强行用Combiner做部分聚合会导致最终结果不确定。如果以后数据量真的到了必须预聚合才能跑完的地步,更合理的做法是改成两阶段任务,先加盐预聚合一次,再去盐做最终合并。
3.4 Driver:作业配置与运行参数
Driver是任务的入口,负责配置作业参数、设置输入输出路径、提交任务并等待结果。这里面有几个参数是有讲究的。
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class MovieDataCleanJob { public static void main(String[] args) throws Exception { if (args.length != 2) { System.err.println("Usage: MovieDataCleanJob <inputPath> <outputPath>"); System.exit(-1); } Configuration conf = new Configuration(); conf.set("mapreduce.output.textoutputformat.separator", "\t"); conf.set("mapreduce.map.output.compress", "true"); conf.set("mapreduce.map.output.compress.codec", "org.apache.hadoop.io.compress.SnappyCodec"); Job job = Job.getInstance(conf, "movie-boxoffice-clean"); job.setJarByClass(MovieDataCleanJob.class); job.setMapperClass(CleanMapper.class); job.setReducerClass(CleanReducer.class); job.setMapOutputKeyClass(Text.class); job.setMapOutputValueClass(MovieWritable.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); job.setNumReduceTasks(4); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); boolean success = job.waitForCompletion(true); if (success) { org.apache.hadoop.mapreduce.Counters counters = job.getCounters(); System.out.println("===== Clean Statistics ====="); System.out.println("RAW_LINES: " + counters.findCounter("CleanStat", "RAW_LINES").getValue()); System.out.println("PARSE_ERROR: " + counters.findCounter("CleanStat", "PARSE_ERROR").getValue()); System.out.println("LESS_FIELDS: " + counters.findCounter("CleanStat", "LESS_FIELDS").getValue()); System.out.println("NULL_KEY_FIELDS: " + counters.findCounter("CleanStat", "NULL_KEY_FIELDS").getValue()); System.out.println("DATE_FORMAT_ERROR: " + counters.findCounter("CleanStat", "DATE_FORMAT_ERROR").getValue()); System.out.println("BOX_OFFICE_PARSE_ERROR: " + counters.findCounter("CleanStat", "BOX_OFFICE_PARSE_ERROR").getValue()); System.out.println("OUTPUT_RECORDS: " + counters.findCounter("CleanResult", "OUTPUT_RECORDS").getValue()); System.out.println("REPAIRED_RECORDS: " + counters.findCounter("CleanResult", "REPAIRED_RECORDS").getValue()); } System.exit(success ? 0 : 1); } }第一,map输出压缩我开了Snappy。这个设置对中间结果特别大、Map到Reduce数据传输耗时长的场景很有效。虽然压缩带来少量CPU开销,但能明显减少磁盘IO和网络传输,实测下来整体任务耗时能减少20%左右。
第二,setNumReduceTasks(4)是综合考虑了数据量和集群资源后的选择。这次数据量30多万行,清洗后约28万行,4个Reduce足以均匀分摊压力。如果后续数据规模翻几倍,这个数字要按数据量动态调整,不能拍脑袋设一个值用到底。
第三,输出分隔符设置成了\t,是为了方便后续用Hive建表加载。默认的分隔符是Tab,但为了防止和字段里的逗号冲突,显式设置一遍更稳妥。
4. 运行流程与效果验证
4.1 数据上HDFS与作业提交
写好的代码打包成jar包,接下来就是标准的Hadoop作业提交流程。先把原始CSV上传到HDFS,再提交任务。
# 创建输入目录并上传原始数据 hdfs dfs -mkdir -p /data/movie/raw hdfs dfs -put movie_boxoffice_2023.csv /data/movie/raw/ # 清理可能存在的旧输出目录 hdfs dfs -rm -r /data/movie/clean # 提交MapReduce作业 hadoop jar movie-clean-job.jar MovieDataCleanJob \ /data/movie/raw /data/movie/clean这里有个习惯性操作:每次跑之前先hdfs dfs -rm -r删除旧输出目录。Hadoop对输出目录的要求是必须不存在,否则直接报错“Output directory ... already exists”。这个坑我一开始踩过好几次,后来干脆写进提交命令里,每次必删。
任务跑完后会有大量日志输出,其中最关键的是最后一段Counter汇总。我通过自定义的CleanStat和CleanResult两组计数器,把原始行数、解析失败数、字段不足数、日期格式异常数、票房解析失败数、最终输出记录数、修复记录数全部打印出来。这些数字就是这轮清洗的质量报告,也是后面验证效果的核心依据。
4.2 清洗前后的数据对比
任务结束后,用下面命令快速查看清洗结果:
hdfs dfs -cat /data/movie/clean/part-r-00000 | head -n 20 hdfs dfs -text /data/movie/clean/part-r-000* | wc -l我这次拿到的样本数据,清洗前的统计结果大概是这样的:
| 指标 | 数量 |
|---|---|
| 原始行数 | 347,852 |
| 字段数不足被过滤 | 3,126 |
| 解析错误被过滤 | 518 |
| 关键字段为空被过滤 | 412 |
| 清洗后输出记录数 | 312,451 |
| 其中修复记录数 | 18,623 |
清洗后,之前那种“40.29亿”和“454400万”混在一起的票房字段,全部统一成了以万元为单位的数值,日期都变成了yyyy-MM-dd格式,评分人数里的千分位逗号全部去掉,类型字段统一用“/”连接,重复记录按电影维度合并。数据从“能看”变成了“能用”。
看几条实际输出:
MOV_001 流浪地球2 科幻/冒险 中国内地 2023-01-22 402900.0 8.3 1023430 郭帆 吴京 repaired MOV_002 满江红 喜剧/悬疑 中国内地 2023-01-22 454400.0 7.0 856120 张艺谋 沈腾,易烊千玺 repaired MOV_003 我不是药神 剧情 中国内地 2018-07-05 30.7 9.0 1617898 文牧野 徐峥 normal每条记录末尾还保留了一个source_status标记,能看出来这条记录是原本就正常还是经过了修复。这个字段在后续数据质量追溯时非常好用。
4.3 数据质量报告与Counters
很多人写MapReduce从来不用Counters,认为它只是任务进度监控的小工具。实际上,Counters是清洗任务里最划算的质量检测手段,比单独写一条统计SQL要省事得多。
我在Mapper里每处理完一类脏数据,就递增对应的Counter。任务结束后,直接从计数器里读出各类问题的数量,一份数据质量报告就自动生成了。比如这次清洗结果显示DATE_FORMAT_ERROR有几千条,说明原始数据的日期格式混乱问题确实严重;BOX_OFFICE_PARSE_ERROR如果突然暴增,那就要怀疑是不是数据源格式变了,而不是代码出了问题。
如果在任务结束后想更精细地看清洗结果,还可以把结果导入Hive做一轮SQL校验。建表语句大致是这样:
CREATE EXTERNAL TABLE movie_clean ( movie_id STRING, movie_name STRING, type STRING, country STRING, release_date STRING, box_office DOUBLE, rating DOUBLE, rating_people BIGINT, director STRING, actor STRING, source_status STRING ) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' LOCATION '/data/movie/clean';然后跑几条简单的校验SQL:
SELECT COUNT(*) AS total, COUNT(DISTINCT movie_id) AS distinct_movie, SUM(CASE WHEN box_office <= 0 THEN 1 ELSE 0 END) AS invalid_box FROM movie_clean;如果total和distinct_movie差距过大,说明去重逻辑还有问题;如果invalid_box很多,就得回去看票房解析是不是漏了某种单位格式。这里也体现了清洗结果和统计口径之间“牵一发动全身”的关系。
5. 常见问题与排查实录
5.1 中文乱码:GBK与UTF-8的恩怨
Hadoop的Text默认按UTF-8解码,但很多影视数据源导出的是GBK编码。直接读GBK文件会出现大量乱码,且这种乱码在日志里不会报错,只会让输出变得不可读。我这次也碰到了,原始CSV用GBK编码,Mapper读进来后片名变成了“娴犳祦鍦扮悆2”。
排查方法很简单:先用file命令确认文件编码,再决定是否需要在Mapper里指定编码读取。如果确定是GBK,一种做法是在读取InputFormat时把整行字节拿出来用InputStreamReader按GBK解码;另一种做法是先统一转码再上传HDFS。我这次为了不改变上游数据链路,选择了在Mapper里做编码兼容。
// 在setup方法中根据参数决定编码 private String inputEncoding = "UTF-8"; @Override protected void setup(Context context) { Configuration conf = context.getConfiguration(); inputEncoding = conf.get("clean.input.encoding", "UTF-8"); } // map中读取时按指定编码解码 // 注意:实际中如果需要完全支持,需要自定义InputFormat, // 但在大多数场景下先转换文件编码是更省事的方案。最省心的还是先本地转码再上传:
iconv -f GBK -t UTF-8 movie_boxoffice_2023.csv > movie_boxoffice_2023_utf8.csv日常建议:在数据管道里把“编码规范化”作为第一道工序,所有进入HDFS的文件统一UTF-8。这样下游所有任务都不用再关心编码问题。
5.2 CSV引号与字段错位
CSV文件里最常见的坑是字段本身包含逗号,比如演员列是“吴京,刘德华”。如果原始数据用了引号包裹整列,而你的解析器没有处理引号,那么一行数据会被切成比预期更多的字段,导致后续字段整体错位。
我在清洗规则里专门引入了CsvParseUtil处理引号,但还有一个更隐蔽的问题:字段数不足或超长。我的处理是:字段数不足直接过滤并计数,字段数超长则尝试拼接,比如把“导演”和“主演”中间多出来的部分合并进演员列。这些规则在代码里就是几个if判断,但实际效果很好。
有一个经验是:不要只依赖字段数量判断错位。有时候字段数正好,但内容已经串列了,比如年份被写到了类型列。这时候要靠字段本身的格式特征来校验,比如检查year字段是否是四位数、类型字段里是否出现“年”字等。清洗规则设计得越贴近业务,效果越好。
5.3 数据倾斜与Reducer热点
清洗任务按电影ID做去重,理论上是均匀分布的,但现实数据里“热门电影”的记录数量可能远超平均。比如某部大热影片在上映期间被多个渠道反复上报,同一个movie_id在map端会被打出几十条记录,这时候负责处理该Key的Reducer就成了热点。
我这次Reducer的负载还算均匀,没有出现明显的长尾。但如果遇到数据倾斜,有几种方案可以尝试:第一种是给Key加盐,比如在Key后面拼一个随机数,让同一条电影的多条记录先分散到多个Reducer做部分合并,然后通过第二轮任务再做最终合并;第二种是采用自定义Partitioner,把预估的热点Key单独分到一个Reducer;第三种是直接提高Reduce任务数量,分散热点压力。三种方案各有适用场景,需要对数据分布有预判后选择。
要注意的是,如果为了性能加了盐,那么清洗逻辑就得改成两阶段MapReduce,代码复杂度会明显上升。数据量没到百万级以上时,我建议保持单阶段简单模型,先把正确性做扎实。
5.4 误删有效数据的风险管控
清洗任务最怕的不是“没清干净”,而是“把有效数据误删了”。我在设计时专门引入了source_status标记,即使某条记录字段异常,也尽量保留并标注,而不是直接丢弃。
比如票房字段解析失败时,我不抛异常,而是置为0.0并递增计数器,同时把source_status置为“repaired”。后续如果想看这些异常记录的原始样子,可以在Reducer里把它们输出到一个单独的“suspicious”目录,供人工复核。
这种“保留+标记”的思路,在数据量不大、可以人工兜底的场景下特别适用。如果数据量太大、人工复核成本过高,也可以退而求其次:把异常记录单独打到一面,定期抽样检查,而不是全量修整。
6. 清洗之后还能做什么
6.1 Hive与Pandas互补的二次校验
MapReduce清洗完的数据,大部分时候还要经过一轮交互式校验才能信任。我的习惯是先把结果加载到Hive,跑几个聚合查询和原始数据的统计结果做交叉验证。比如手动算一下总票房是否符合预期、不同类型电影的评分均值是否合理、2023年上映电影数量是否和公开资料对得上。这些校验不复杂,但能发现清洗规则里的逻辑漏洞。
如果需要更灵活的探索性分析,我会用Pandas再对清洗后的文件做一轮可视化前的预处理。这个阶段pandas的灵活性就体现出来了,可以快速画分布直方图、做数据透视。MapReduce负责“大规模清洗+结构化落地”,Pandas负责“小规模校验+探索性分析”,两者配合是我目前最顺手的工作流。
6.2 票房分析的下游应用
清洗完成的数据可以支撑很多下游分析任务:按年份统计票房趋势、按类型分析观众偏好、计算导演和演员的票房号召力、分析评分和票房的相关性等。我这次清洗完之后,用Hive跑了一个“历年票房Top20”的统计,结果比清洗前靠谱太多了——之前因为单位不统一,一部40亿票房的电影被算成40.29,差点排进倒数。
所以数据清洗不是终点,它是所有后续分析的基础设施。清洗做得好不好,决定了上层分析的结论可不可信。很多人花时间调模型、调参数,却忽视了数据质量这个最根本的问题,这是我觉得最值得提醒的一点。
6.3 把这个任务做进调度
清洗任务如果只跑一次,手动提交就够了。但如果数据每天更新,就需要把整个流程配置到调度系统里。我这边用的是简单的crontab加Shell脚本,每天凌晨把当天新增的CSV导入HDFS,然后提交清洗作业,结束后把结果表分区写入Hive。MapReduce任务的幂等性很好,输入不变输出不变,所以调度重跑也不会造成重复写入的问题。
如果团队已经有Airflow或Oozie这类调度平台,把清洗任务封装成一个可调度的节点是更规范的做法。不过无论用哪种调度方式,核心思路都是一样的:数据必须先过清洗这关,才能进入下游的分析和展示链路。