ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

数据清洗在大数据中的作用与前沿技术实践

数据清洗在大数据中的作用与前沿技术实践 做了这么多年的数据项目我越来越觉得数据清洗才是真正决定项目成败的环节。很多人一提到大数据就想到Hadoop、Spark、实时计算这些闪闪发光的词却往往忽略了一个朴素的事实数据进到系统里的那一刻是脏的后面所有计算、分析、建模都建立在数据到底可不可信这个问题之上。数据清洗本质上就是在大数据链条上守住第一道可信关口的活儿。这篇文章我想结合这些年在不同项目里摸爬滚打的经验聊聊数据清洗技术在大数据领域的前沿动态——从工具演进、技术选型、实操细节到一些面试常见考察点尽量写清楚哪些做法值得参考哪些坑真的不值得再踩一遍。适合正在做数据开发、算法工程或者想系统梳理数据清洗知识体系的朋友。1. 数据清洗在大数据链路中的位置与核心挑战1.1 为什么数据清洗越来越被单独当成一个技术方向早些年做数据清洗其实就是写几个SQL把空值过滤掉、把重复记录去一下感觉这活儿是个体力活谁都能干。但这几年行业风向明显变了。数据湖、数据仓库、实时数仓、大模型训练这些场景铺开之后数据清洗从一个辅助动作慢慢变成了一个需要独立设计、独立维护的技术模块。我自己感受最深的是三件事。第一数据的规模已经不是跑个脚本处理一下就能覆盖的。以前一张表几百万行Excel或者pandas读进来洗一洗就完了。现在一张表几亿行甚至几十亿行清洗逻辑稍微复杂一点跑一次就要几十分钟甚至几个小时。规模上来了原先怎么洗的问题就退居其次在什么架构上洗按什么频率洗怎么验证洗没洗干净反而成了核心问题。第二数据源比以前杂得多。业务库、日志系统、埋点数据、第三方接口、物联网设备数据全都往里灌。每个源的字段命名不一样、时间格式不统一、维度值的取值逻辑不同。以前数据清洗面对的是同一张表里的脏数据现在面对的是不同系统之间的脏数据清洗的理解维度从记录级扩展到了模型级。第三大家对清洗结果的要求变了。以前清洗完能出数就行现在上层往往挂着实时报表、用户画像、风控模型。清洗质量直接影响业务决策一个字段填错就可能让一个策略误判。所以清洗这条线开始有了自己的质量指标、监控告警、甚至独立的数据质量团队。这也是为什么现在招聘市场上数据质量管理数据治理这类岗位越来越多——数据清洗的技术含量确实被重新定价了。1.2 数据清洗面临的三类典型挑战我把实际项目里最头疼的问题归纳成三类规模化、多样性、时效性。规模化这块核心矛盾是数据量越来越大清洗逻辑越来越复杂之间的矛盾。举个典型的例子做用户画像的时候需要对几亿用户的经纬度做逆地理编码同时还要关联数十个维度表做标签修正。这种量级的清洗任务单机已经没戏了必须考虑分布式执行。我之前接过一个项目日志数据一天就上百亿条清洗规则还涉及跨天去重、会话拼接直接写MapReduce跑不动后来用Spark按用户ID分桶、按会话窗口聚合才把耗时从十几个小时压到两个小时以内。多样性这块最难的不是数据本身脏而是脏的方式千奇百怪。同一个用户ID在A系统里叫user_id在B系统里叫uid在C系统里还带前后缀。手机号有的带86有的不带有的中间有空格。地址字段更是重灾区同一个小区能被写出七八种叫法。这些问题的本质是缺少统一的数据标准清洗时要做的其实不只是修数据而是通过映射、归一化、字典表这些手段把多源数据拉到同一个语义平面上来。时效性这块和批处理完全不一样。实时数仓和实时风控场景下数据清洗要在流上完成延迟要求是秒级甚至毫秒级。批处理里那些全表扫描、排序去重、回填历史的操作在流式环境下很多都不适用得改成窗口计算、状态管理、增量清洗的思路。我们内部有句玩笑话批处理清洗是洗完再穿流式清洗是边穿边洗两者的工程复杂度完全不在一个量级。这三类挑战不是孤立的真实项目里往往同时出现。所以我一直觉得数据清洗不是一个静态的技能点而是需要跟着数据架构走的一套方法论。这也是这篇文章想传递的核心视角。2. 工具链演进从单机脚本到分布式质量体系2.1 pandas依然是主力但有明显的边界先聊pandas。无论外面出了多少新工具pandas在数据清洗这块的地位短期还是很难被替代。它语法简单、生态成熟、处理中小规模数据非常顺手。我在日常工作中凡是能落到本地或者单机处理的清洗任务第一反应还是pandas。不过pandas有一个必须清醒认识的边界它默认会把数据全量加载到内存里。一张几千万行的DataFrame加上几列文本字段内存占用轻松超过20GB。这也是为什么很多人会拿Qt 表格大数据卡顿优化来类比——本质上都是同一个问题数据量一旦上去全量处理的方案就会在性能上卡脖子。Qt表格那个场景里从QTableWidget换到QTableView加自定义QAbstractTableModel核心思路是视图只渲染可见的那几十行数据留在模型里按需取用pandas的破局思路也很像就是分块、采样、避免一次性全量展开。有人可能会问那直接在pandas里把所有数据读进来不就行了能行的时候当然行但一旦接近内存上限各种奇奇怪怪的问题就来了。比如内存碎片导致OOM、字符串类型占用过高、value_counts在超大数据集上卡住。实际操作里我一般会先用pd.read_csv的参数做预处理比如usecols只读需要的列、dtype指定列类型、nrows先抽样探路把数据压缩到合理范围再进内存。再补充一个pandas的实用小技巧能用向量化操作就别用apply。很多人在清洗时习惯写个自定义函数然后df.apply一下子数据量小无所谓数据量一大就是灾难。向量化操作底层是C循环apply是Python循环两者性能差距可以到几十倍甚至上百倍。真到了必须逐行处理的时候要么改用numba要么考虑把活儿交给分布式引擎。2.2 Spark与Flink分布式清洗的两种典型形态当数据规模超出单机极限就要上分布式。目前最主流的两类引擎就是Spark和Flink。Spark做批式清洗核心是DataFrame API加Spark SQL。我一般会把清洗逻辑拆成一个个Spark任务挂到集群上跑。这里想强调一个经验清洗任务非常吃Shuffle尤其涉及join操作的时候数据倾斜是最大的敌人。比如把几亿条明细表和一张维度表join如果维度表不大应该优先考虑广播变量避免每个task都去拉全量维度数据。这个细节能直接影响任务能否按时跑完。集群部署策略上如果只是做清洗没必要上太重的资源driver内存给足、executor数量按数据量估算即可。我们的经验是先做数据抽样估算体积再按经验公式来配资源。举个例子一个每天跑一次的离线清洗任务数据量在500GB到1TB之间我会申请10到20个executor每个4核8GBdriver给4GB。如果清洗逻辑里有大量的groupBy和joinspark.sql.shuffle.partitions要适当调大否则大量数据堆在少数task上一样会拖垮任务。如果只是简单的过滤和格式转换executor数量可以保守一些。另外存储上尽量按日期分区清洗任务只扫描当天的分区能省下大量IO开销。Flink则负责流式的数据清洗。在实时场景里数据是源源不断进来的没法像批处理那样等全量再动手。Flink的清洗逻辑通常写成DataStream算子链过滤、转换、窗口去重、状态管理。最常见的一个需求是实时去重比如清洗埋点日志时去掉重复点击。批处理里用distinct或者group by就行但流式里得用状态存储加TTL把一定时间窗口内的主键记录下来超过窗口的自然遗忘。这种思路和批处理差异很大但对实时数据质量至关重要。Spark和Flink的选型其实不冲突。我们的架构里离线链路用Spark做每日清洗实时链路用Flink做秒级清洗两套引擎共用一个数据质量层。业务需要的是今天的数据最迟到明天早上是干净的——Spark做得到需要的是当前这一秒的数据不能有重复——那就得上Flink。2.3 数据质量框架把清洗规则沉淀成工程资产前几年做数据质量基本靠人肉。每个清洗任务跑完人工抽查几条看看有没有为空、有没有明显异常然后拍脑袋说感觉还行。现在这种做法越来越不够用了因为数据规模和业务复杂度上来后人肉抽查根本覆盖不到边缘case。所以行业内开始流行数据质量框架帮我们把规则沉淀成工程化资产。目前用得比较多的有Great Expectations、Deequ亚马逊开源的基于Spark的库、还有各类商业方案。以Deequ为例它的核心思路是把数据质量断言写到代码里让测试框架来校验数据。比如user_id列不能为空金额字段必须大于0订单日期必须在合理范围内这些规则都可以转换成自动化的断言。每次清洗任务跑完之后自动跑一轮断言不通过就发告警。这样做最大的好处是清洗质量的问题从事后发现变成了事前拦截。Great Expectations则更像是数据文档加校验工具的结合体。它可以生成数据的profile报告把你的期望值写成expectation每次数据更新后自动验证。我比较推荐在数据湖的场景里用因为它能帮你维护一套关于数据的约定新来的同学看到期望文档很快就能理解每张表、每个字段的质量标准。需要强调的是数据质量框架不是银弹。它解决的是如何验证清洗结果无法替代如何设计清洗逻辑。框架用得好前提是想清楚数据质量标准是什么、业务的容忍阈值是多少。这些还是要人来定。2.4 前沿方向从规则清洗走向智能清洗传统的清洗都是基于明确的规则if 字段为空 then 填充if 值超出阈值 then 标记。这种规则驱动的清洗方式成熟可靠但有一个问题规则需要人来写而写规则的前提是你已经知道脏数据长什么样。现实是很多数据源的脏数据形态超出预期规则写起来没完没了。所以这两年行业里开始出现智能清洗的探索。大致有几个方向。一个是异常检测的机器学习化。以前用IQR、Z-Score这些统计方法现在越来越多用孤立森林、自动编码器等无监督模型来识别异常模式。尤其是面对高维数据时单字段的统计阈值检测不到组合异常比如一个用户IP来自国外但收货地址在国内且下单时间在凌晨这种多维异常靠传统方法很难发现而模型可以同时捕捉多个特征的偏离。另一个方向是大模型辅助清洗。大语言模型在理解自然语言和上下文方面有天然优势可以用来做规则自动生成、字段语义识别、甚至直接处理非结构化文本清洗。比如给模型一段样本数据让它推断出可能的清洗规则再由人工确认后固化成代码。这条路还在早期落地的不多但确实是数据清洗领域值得关注的前沿信号。还有一个方向是清洗逻辑的自动化编排和DataOps。把数据探查、清洗、校验、告警全部编排成可配置的流水线降低人工介入的频率。我在一些成熟团队看到过类似平台清洗规则可以在界面上配置数据质量报告自动生成异常自动触发重新对齐流程。这已经不是单纯的脚本工程而是一个数据产品了。这些方向对实际工作的启示是数据清洗不会被淘汰但会被重新定义。做数据开发的人如果只停留在会用pandas去重的层面确实可能被工具取代但如果能理解数据质量和业务口径有什么关系异常数据背后的业务含义是什么反而会因为自动化程度的提高而更能发挥人的判断价值。3. 实操从原始数据到可用数据的完整清洗流程3.1 第一步数据探查别急着动手我见过很多新人拿到数据就开始写清洗脚本结果跑了半天发现字段理解错了、单位没换算、数据早就被上游规则处理过。所以每次做清洗之前我强烈建议先做数据探查。这步花的时间越长后面返工的时间越短。数据探查说白了就是搞明白三件事数据长什么样、数据质量如何、数据和业务期望之间有没有偏差。以pandas为例最基本的探查就是df.info()和df.describe()前者看字段类型和缺失情况后者看数值分布的统计量。不过这两个方法在面对复杂数据时远远不够我一般还会用以下手段对每个维度字段做value_counts()看枚举值的分布是否合理有没有奇怪的取值对时间字段做min/max检查确认时间范围是否符合预期有没有未来时间或者年份异常的记录用isnull().sum()统计每个字段的缺失量再按业务维度分组看缺失的分布规律对数值字段做分位数检查特别关注p99、p999判断有没有极端值这些操作看起来简单但能提前暴露大部分问题。我的习惯是探查完之后要先写一份简短的数据质量摘要把发现的问题列出来再和业务方确认哪些是正常的比如空值本身代表合法语义哪些需要清洗。这里提醒一句不要拿到数据就直接按教科书处理一定要先问自己这个空值在这里是不是真的有业务含义。3.2 缺失值、异常值与重复值的处理细节探查清楚了就可以动手清洗。缺失值处理上我的原则是能保留就不乱删能填补就不留空。删除确实是最省事的办法但代价是损失信息尤其当缺失比例比较低时影响不大但缺失比例超过30%的时候删除可能就把整个样本的分布带偏了。填补的方式要看语义连续型数值字段可以用均值、中位数填补我个人更推荐中位数因为它对异常值不敏感时间序列数据用前向填充ffill往往比全局均值更合理分类字段可以填众数或者新增一个未知类别。还有一种是模型预测填补用已有字段构建模型去预测缺失值效果好但成本也高只有在缺失值很关键的时候才值得做。异常值检测上最简单常用的是IQR规则和Z-Score。IQR规则比较稳健适合分布不那么规整的数据。比如一个字段的25%分位数是Q175%分位数是Q3超出Q1-1.5IQR或者Q31.5IQR的就可以标记为异常。Z-Score适合近似正态分布的字段。但这里要注意异常值不等于错误值有时候异常值恰恰是业务想找的极端情况。比如反欺诈场景里金额特别大的交易不是要清洗掉而是应该标记出来作为重点分析对象。所以清洗异常值之前一定要想清楚这个异常的判定是业务规则还是统计学规则。重复值处理也有讲究。最简单的df.drop_duplicates()只能处理完全重复的行但真实场景里很多重复是不完全重复——比如同一个人下了两笔订单只有订单号不同其他字段都一样。这种要靠业务主键来判断。另外基于时间窗口的重复比如同一天内多次点击可能不是数据错误而是正常的用户行为去重前必须明确业务语义。下面给一段比较完整的pandas清洗示例代码import pandas as pd import numpy as np # 1. 读取数据只保留需要的列并指定类型减少内存 df pd.read_csv( raw_data.csv, usecols[user_id, order_id, amount, status, order_time], dtype{user_id: str, order_id: str}, parse_dates[order_time] ) # 2. 缺失值处理金额缺失用中位数填充状态缺失填UNKNOWN df[amount] df[amount].fillna(df[amount].median()) df[status] df[status].fillna(UNKNOWN) # 3. 异常值检测金额超过p993*IQR的标记为异常但不直接删除 q1 df[amount].quantile(0.25) q3 df[amount].quantile(0.75) iqr q3 - q1 upper_bound q3 3 * iqr df[amount_anomaly] df[amount] upper_bound # 4. 重复值处理按业务主键去重保留最早一条 df df.sort_values(order_time) df df.drop_duplicates(subset[user_id, order_id], keepfirst) # 5. 标准化状态字段统一为小写 df[status] df[status].str.lower() print(df.shape) print(df[amount_anomaly].value_counts())这段代码覆盖了最常见的清洗动作。实际项目里中间每一步都需要补充日志、统计清洗前后数据量的变化便于追溯。注意凡是涉及是否删除的清洗决策都建议先标记而不是直接删。真实项目里数据被删掉往往就再也找不回来了标记异常字段至少保留了追溯和恢复的可能。3.3 数据标准化与一致性校验接下来是标准化。标准化解决的是同一个意思、多种表达的问题。几个高频场景时间字段。不同系统可能给到2024-01-01 00:00:002024/01/0100这些完全不同的格式。统一处理时我一般直接转成标准时间格式再存比如都用ISO 8601。同时要注意时区问题大数据场景里上下游系统常跨时区清洗时如果不统一后面做时间窗口计算或者报表统计就会错位。编码问题。早期系统常有编码混乱的情况中文字段可能混着GBK、UTF-8。清洗时要么统一转换成UTF-8要么在读取时指定正确的编码。遇到乱码不要上来就替换掉先判断是转码错误还是真的垃圾字符能找到原始编码的用iconv转回来找不回的直接剔除。单位和量纲。这一点最容易被忽视。比如金额字段有的系统存分、有的存元重量字段有的存克、有的存千克。清洗时务必统一单位否则聚合统计就是灾难。我的做法是在表结构设计阶段就定好每个字段的标准单位清洗任务强制转换。维度值的归一化。同样是性别A系统存男/女B系统存1/0C系统存M/F。清洗时通过映射字典统一。这里有个经验映射关系一定要做成配置表放到仓库里别埋在代码里否则后面维护的人会疯掉。标准化做完之后最后一步是跨表的一致性校验。比如用户表和订单表的user_id要对得上订单里的城市代码要在城市维度表里存在。校验不通过的数据要么修复、要么拦截。我自己用Spark做这部分时会比较喜欢用left_anti join找出那些订单关联不上用户的记录单独落一张异常表留着给业务方分析原因而不是直接删掉。数据清洗的很多决策都需要留下痕迹。4. 常见问题与排查技巧实录4.1 内存爆炸与性能劣化写到这里必须分享几个踩坑经历。第一个就是内存爆炸。曾经处理一个十几GB的日志文件我用pandas直接pd.read_csv()加载结果还没读到一半内存就爆了。后来改用分块读取每块50万行清洗完再合并才把问题解决。分块时要小心的是如果清洗逻辑里有全局统计比如要计算全表的中位数必须先扫一遍拿到全局统计量再回到分块数据里做填充。这其实和Qt表格大数据卡顿优化里视图只渲染可见行但模型要维护全量数据索引的思路是一致的——全局信息不能丢但展示和处理可以分片进行。另一个性能坑是字符串类型的隐形炸弹。pandas里有些列本应是category类型读进来却成了object每个字符串都是独立的Python对象内存占用大得惊人。用df.astype(category)把维度字段转成类别类型内存能省好几个量级。大批量处理时这块优化效果立竿见影。还有一种常见问题是内存泄漏。在循环里反复调用pandas操作、每次生成新DataFrame却没有释放旧的引用程序跑着跑着内存就涨上去了。排查时用gc.collect()不一定管用更可靠的是把循环体拆出来单测看单轮处理大概占多少内存再估算整个任务的峰值。必要时用tracemalloc定位分配热点。4.2 规则失效数据漂移与血缘第二个大坑是数据漂移。所谓漂移就是上游数据格式或者分布悄悄变了但你写的清洗规则还停留在上次的状态。比如上游某天把金额字段从元改成了分你的清洗逻辑没有感知出来的数据直接大了一百倍。这类问题在纯写逻辑时很难发现所以我现在的习惯是清洗任务跑完后必须自动做一轮数量级校验比如金额字段的均值如果比昨天翻了十倍就必须告警暂停。这其实也是一种数据质量断言只是比较粗糙但非常有效。还有一个容易被忽略的点是血缘关系。数据清洗是会不断迭代的同一个指标可能经过了四五层清洗加工。没有血缘追踪的话某个口径出了问题你去追源头会发现根本不知道是哪一层改坏了。现在业界比较好的做法是推进元数据平台采集血缘清洗任务的输入、输出、转换规则都记录在案。没有这个条件的话至少也应该在清洗代码里写清楚注释和版本号——我见过太多这个SQL不知道谁写的但都在跑的恐怖场景。4.3 面试中数据清洗的高频考察最后聊一下大数据领域的面试。数据清洗在面试里的出镜率一直很高而且考察点往往很综合。最常见的面试题是让你完整说一遍清洗流程这时候要重点展现数据探查-规则制定-执行清洗-质量校验-结果验证的闭环思维而不是简单列举去重、填补空值这些动作。另一个高频考点是工具对比。面试官会问pandas和Spark的区别、什么时候选哪个。回答的关键在于理解数据规模和计算模型这两个维度。你需要讲清楚pandas适合单机、中小规模、交互式探索胜在灵活和生态Spark适合分布式、大规模、批处理胜在横向扩展和统一SQL接口Flink适合实时流处理胜在低延迟和状态管理。不要只报菜名最好结合自己实际做过的场景比如我当时处理几亿行数据时pandas内存扛不住换到Spark以后用分区加广播变量把任务跑通了这样比背一堆API名称要有说服力得多。还有一类是场景题比如给你一个几亿行的日志表里面字段有乱码、有重复、有时间区不一致你怎么在大数据平台上做清洗这种题考察的不只是技术更是工程能力。我的解题思路一般是从数据探查入手先说如何抽样、如何用describe和value_counts发现问题然后给出分批读入、分布式处理、质量校验的方案最后强调要留审计日志和告警机制。把思路讲得完整、可落地面试官通常会比较认可。数据清洗的面试题还有一个变体数据质量指标。比如问你如何衡量清洗结果的好与坏。我一般会回答几个维度完整率缺失比例、准确性抽样校验正确率、一致性同义字段取值对齐度、及时性清洗任务是否按时完成。这几个维度可以支撑起一套简要的数据质量评估体系也是日常工作中很有用的框架。另外多说一句数据科学与大数据技术方向的同学在准备面试时不要只背组件和框架的八股文。面试官真正想看到的是你面对一堆不听话的数据时能不能有条理地拆解问题、设计流程、验证结果。数据清洗恰恰是这种能力最好的试金石。我个人在实际操作中体会到数据清洗这件事做到后面拼的不是工具用得有多熟练而是你对业务数据的理解有多深。同样一个空值在不同业务里可能有完全相反的含义同样一个重复记录在日志分析和订单统计里的处理方式也截然不同。所以我建议刚开始接触数据清洗的朋友别急着追求花哨的算法和框架先老老实实把一个数据集摸透从数据探查开始把每一步的取舍理由记录下来。另外一个小技巧是每次清洗任务上线前都先写清楚清洗前多少行、清洗后多少行、每个规则影响了多少数据这些数字能帮你在数据出问题时快速定位。数据清洗是苦活、细活但也是整个大数据链条里回报率最高的投入之一把这个环节做扎实了后面所有分析、建模、决策都会省心不少。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进