尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

数据清洗实战:从Pandas到Spark的完整方法论

数据清洗实战:从Pandas到Spark的完整方法论 干这行这么多年我越来越觉得数据清洗才是大数据项目里真正见功底的地方。很多人以为搭个集群、跑个分析模型就是大数据实际上一套数据从采集端到可视化大屏清洗环节通常要吃掉整个项目60%~80%的工期。招聘数据、农产品价格、网约车轨迹、校园一卡通流水我经手过的项目里几乎没有一次能幸免于脏数据的骚扰。这篇文章我就结合平时做过的pandas清洗、Spark清洗、以及MapReduce场景下的清洗任务把数据清洗这件事从思路到落地完整拆一遍重点讲清楚每个环节为什么这么做、怎么做最省事。1. 数据清洗到底在“洗”什么1.1 脏数据都是从哪里冒出来的先说个我自己的体会脏数据不是某一个环节造成的它是整条生产链路“攒”出来的问题。业务系统里人为填写的字段录入员手滑、下拉框选项过期、字段被复制粘贴错位这是第一层脏。前端埋点漏埋、SDK版本升级导致字段名更换这是第二层脏。第三方接口返回的数据结构不稳定昨天还是JSON数组今天变成嵌套对象这是第三层脏。还有多张表关联时主键口径不一致、不同时区的服务器各自记录时间、爬虫抓下来的数据带上一堆HTML标签和非法字符第七层第八层脏都有。我曾经处理过一批网约车订单数据司机端上报的经纬度有时是WGS84坐标、有时是GCJ-02坐标混在一张表里根本没法直接做距离计算。这类问题靠什么发现靠的是清洗前的数据探查也就是先别急着写清洗代码花时间把每个字段的分布、空值率、格式、取值枚举全部摸一遍。数据探查做得越细后面清洗规则就越有底气。我见过不少新手拿到数据就开始写dropna、去重结果把有用的记录删掉一大半这就是典型的跳过探查直接动手。1.2 大数据场景下清洗要盯住五个质量维度数据质量在不同资料里有不同定义落到实际项目中我习惯归纳成五个维度去检查这五个维度也是清洗规则设计的纲领。维度要解决的问题典型表现完整性字段是否缺失或为空手机号为空、年龄为空、地址缺省唯一性同一条记录是否重复同一订单出现多次、用户重复注册准确性取值是否真实且在合法范围内年龄500岁、金额为负数、经纬度在海里一致性多表、多源数据是否口径统一单位不同、编码不同、时区不同时效性数据是否过期或时间戳错乱订单时间晚于当前时间、日志时间乱序这五个维度不是孤立处理的。举个例子一张订单表里user_id为空这不仅影响完整性还直接影响去重和后续join。所以清洗规则的执行顺序很重要一般先做格式和类型归一再做去重再处理缺失和异常最后做跨表一致性校验。顺序反了会出问题先删重复再填缺失可能把两条不同用户记录合并弄丢先填缺失再格式校验可能填进去的数据类型根本不对。1.3 清洗在大数据链路里的位置大数据架构通常分采集、存储、计算、应用四层清洗横跨存储和计算两层。实时链路里清洗逻辑嵌在Flink作业中流过来一条处理一条离线链路里清洗逻辑体现在ETL脚本、Hive SQL、Spark作业里对ODS层原始数据做处理后落入DWD层。我个人的建议是清洗前置越早越好。原始数据落库之后第一时间做清洗后面分析、建模、可视化用的都是干净数据。如果清洗拖到分析阶段再做每个下游任务都要重复处理脏数据跑的慢不说各任务清洗口径不一致还会导致相同指标算出来两个数。这也是为什么有些团队专门搭数据质量检查框架每天定时对表的数据量、空值率、主键唯一性做监控发现异常立刻报警把脏数据挡在入仓之前。2. 清洗方案的整体设计与工具选型2.1 清洗流程的分层思路做数据清洗不能上来就写代码我习惯把整个流程拆成四层探查层、规则层、执行层、校验层。探查层负责回答“数据现在什么样”包括字段数量、记录数、每字段的非空率、枚举分布、类型推断、重复率、异常值分布。执行方式是写一些统计SQL或者pandas的describe、value_counts。规则层负责回答“怎么把数据洗干净”把探查发现的问题翻译成可执行的清洗规则。比如“age字段存在大于120的值需要置为空并标记”“create_time存在两种格式需要统一成标准时间戳”。规则必须写成文档哪怕项目赶工期也要留个简版不然半年后你看着自己写的清洗脚本都不知道当时为什么删那批数据。执行层就是写代码跑清洗任务单机用pandas分布式用Spark或者MapReduce。执行层要注意幂等性同一份原始数据无论跑多少遍产出的干净表必须完全一致这样重跑任务才不会产生重复的数据变更。校验层是最后一道关清洗前后记录数对比、关键字段空值率对比、抽样人工检查把校验结果写进日志或报表。我自己项目里一定会留一份清洗报告字段级空值率从多少降到多少、去重多少条、异常值修正多少条写清楚。这个报告是后续跟业务方扯数据口径时的底牌。2.2 pandas、Spark、MapReduce怎么选工具选型这问题几乎每个新人都问过我的回答是看数据量、看开发效率、看运行环境。pandas适合数据量在单机内存范围内的场景几GB以内最舒服。它的优势是交互性强写几行代码立刻出结果特别适合做探查和规则验证。我先用pandas在小样本上验证清洗逻辑是否正确再把逻辑迁移到Spark跑全量。这种方式既能快速迭代又能避免全量任务跑完才发现逻辑写错。Spark适合处理GB到PB级的离线数据。DataFrame API比RDD写起来顺手得多处理缺失值、去重、类型转换、字符串清洗这些操作都有现成函数而且Spark SQL能用标准SQL写清洗逻辑团队里会SQL的人直接能上手。现在网约车、电商、校园大数据这类项目里离线清洗基本都是Spark的活。MapReduce在数据清洗里不算主力但两种场景很常见。一是招聘数据清洗这类教学案例里要求用MapReduce实现练的是分而治之的思路二是某些老集群没有Spark只能在MapReduce里写清洗逻辑。MapReduce开发效率低、调试麻烦但它的优势是稳定、可控、不依赖额外组件适合逻辑简单的过滤、格式化、去重任务。维度pandasSparkMapReduce数据量单机内存内GB级分布式TB/PB级分布式吞吐大开发效率高交互式探索较高DataFrame/SQL低需要完整工程代码调试体验即时反馈提交作业看日志日志排查较慢典型场景探查、小样本验证离线批量清洗ETL教学案例、老集群作业学习门槛低中较高2.3 清洗规则怎么沉淀成资产很多团队做完一个项目清洗代码就废弃了这是很可惜的。清洗规则其实是可以沉淀的资产我建议做三件事。第一把规则配置化。比如“字段长度大于X判定为异常”“枚举值只允许集合A、B、C”这些条件不要写死在代码里而是放进配置文件或规则表下次换数据源改配置就行。第二建立字段字典。每个字段标注业务含义、数据类型、取值范围、清洗规则、负责人。字段字典是数据治理的基础也是新同学上手最快的学习材料。第三规则要有版本。清洗规则会随业务变化而调整比如业务新增了一个地区的编码规则就要更新。规则跟代码一起做版本管理每次改动在注释或文档里写明原因出了问题能回溯。3. 核心清洗场景与实操要点3.1 缺失值先判断能不能删再决定怎么填缺失值处理是在“删除”“填充”“保留处理”“不处理”之间做选择不是无脑fillna。我踩过的坑告诉我第一件事永远是问业务方这个字段缺失代表什么如果一份农产品价格数据里“产地”字段缺失可能只是因为录入时没有这个信息删掉整行会损失有效价格记录不删又没法按产地分析。这时候用“未知”填充并单独打标签比删除更合理。填充策略要分情况。数值型字段比如价格、年龄、金额如果分布近似正态用均值填充如果存在明显偏态和离群点用中位数更稳时序数据用前向填充或者线性插值。分类字段用众数或者专门的“未知”值。还有一种是建模推填用其他字段训练模型预测缺失值适合缺失率不高且字段之间强相关的场景但要控制成本。删除要讲究条件。记录数少且缺失字段对分析无用可以整行删除某个字段缺失率超过80%这个字段基本可以放弃。删除前一定要记录删除数量别让下游觉得表凭空少数据。import pandas as pd df pd.read_csv(price_data.csv, encodingutf-8-sig) # 统计空值率决定删除还是填充 null_ratio df.isnull().mean() # 数值字段用中位数填充 df[price] pd.to_numeric(df[price], errorscoerce) df[price] df[price].fillna(df[price].median()) # 分类字段用未知值填充 df[origin] df[origin].fillna(未知) # 缺失率过高的字段直接删除 df df.drop(columns[col for col in df.columns if null_ratio[col] 0.8])3.2 去重完全重复好去业务重复难去去重分两层。第一层是整行完全重复这种最简单drop_duplicates一行搞定。第二层是业务意义上的重复两条记录字面上不完全一样但表示的同一个业务事件。比如同一订单id出现两次但运费字段一个有值一个没值同一个用户id注册了两次邮箱不同。这种去重必须指定业务主键再决定保留哪条。保留策略我最常用的是两种按时间保留最新一条或者按某字段的填充程度保留信息更全的一条。可以用sort_values排序后再去重也可以用groupby配合取第一条。还要注意去重不能只盯一张表跨表的主键去重在数仓里通常用row_number窗口函数实现处理逻辑更清晰。# 完全重复 df df.drop_duplicates() # 业务主键去重同一订单保留最新时间的一条 df df.sort_values(create_time, ascendingFalse) df df.drop_duplicates(subset[order_id], keepfirst)Spark里同样有dropDuplicates方法指定subset参数即可。但分布式中要注意一个坑去重发生在shuffle阶段如果数据量极大且主键重复率很高容易引发数据倾斜个别executor数据量暴增拖垮整个任务。这种情况可以给主键加随机后缀先分散再去掉后缀做最终去重不过绝大多数清洗场景用不上这个技巧知道有这回事就行。3.3 格式归一类型、编码、时间、地理坐标格式归一的目标是让每个字段就一种格式这是后续计算的基础。我举几个高频场景。时间字段是最乱的常见的有“2024/01/05 08:30:00”“2024-01-05 08:30”“2024年1月5日”混在一起。统一用to_datetime解析解析不了的置为NaT再单独处理。字符串类型的手机号、身份证号看起来是数字但不能当数值处理必须在读入时就指定dtypestr防止前导零丢失。地理坐标也是个经典坑。GPS设备、地图SDK、第三方接口各自输出不同坐标系直接按数值计算距离会产生几公里甚至几十公里的误差。清洗时要根据数据来源识别坐标系类型统一转换后再落库。这个转换逻辑很成熟网上有现成库但规则层必须写明每批数据的坐标系来源。正则表达式是做格式归一最顺手的工具。手机号、邮箱、邮编的校验都靠正则。跑大数据集群之前先用正则在小样本上跑一遍确认没有误伤再上线全量任务。我见过一个案例清洗规则把订单备注里的“电话138xxxx8888”全部误判为非法手机号差点把备注信息清空就是正则边界没写好。import re def clean_phone(value): text str(value).strip() if re.fullmatch(r1[3-9]\d{9}, text): return text return None df[phone] df[phone].apply(clean_phone)3.4 异常值用业务阈值为主统计方法为辅异常值处理最容易走极端。一种是把所有超出均值加减三倍标准差的值全部删除结果把真实的高消费用户当成异常清掉了另一种是看到异常值就手工改改完没有记录后面解释不清。我的原则是业务规则优先于统计规则。业务规则就是字段本身的合法范围。年龄在0到120之间、折扣率在0到1之间、订单金额大于0这些边界直接判非法。统计方法适合没有明确业务边界的连续字段比如商品价格、配送时长用分位数或者Z-score做辅助判断。用IQR方法时上下界以外的点不一定要删除可以标记出来让业务方确认毕竟离群点有可能代表异常事件或者高价值用户。处理动作上非法值转缺失再按缺失策略处理是干净的做法。标记异常值也很重要给数据加一列quality_flag值为normal、missing、abnormal让下游能追踪每一行数据的清洗痕迹。增加一个标记字段比直接物理删数据稳妥得多这也是数据治理里的可追踪性原则。# 数值转换非法值转NaN df[age] pd.to_numeric(df[age], errorscoerce) # 业务规则年龄合法范围 df.loc[(df[age] 0) | (df[age] 120), age] None # 数值字段异常值用IQR识别并打标 q1 df[amount].quantile(0.25) q3 df[amount].quantile(0.75) iqr q3 - q1 lower, upper q1 - 1.5 * iqr, q3 1.5 * iqr df[amount_flag] df[amount].apply( lambda x: normal if lower x upper else abnormal )3.5 多源一致性编码映射和时间口径统一大数据项目几乎都是多数据源汇合每个源来一套自己的编码这是常态。做得比较多的两件事是维度编码映射和时间口径统一。维度编码映射比如省份编码一个源用“110000”另一个源用“11”还有一个源直接存“北京市”。清洗时要建一张映射表把各源编码统一到标准编码上。在Spark里做这个映射一是把小维表转成字典后用udf处理二是直接与小表做join我推荐join方式性能稳定。时间口径统一更隐蔽。订单表存的是北京时间日志埋点存的是UTC时间日志分析要按天聚合时就差出8个小时。清洗规则必须明确统一时间基准通常全链路统一用UTC存储、展示时再转本地时区这样可以避免不同地区服务器各自为政。还有夏令时地区的数据源时间不一致问题更严重抓取第三方数据前先把时区规则问清楚。这些一致性工作看似不起眼却决定了多表关联时能否对得上。我见过团队做用户留存分析两张表的日期口径差8小时算出留存率明显异常排查了两天才发现是时区问题。所以多源数据的口径说明文档一定要在清洗前写好而不是等对不上数再回头补。4. 实操过程一套从pandas到Spark的完整清洗链路4.1 第一步用pandas做样本探查和规则验证拿一个招聘数据清洗的案例举例这个案例很多教材里都有实际项目中基本也是这个思路。原始数据大概是这样的一条岗位记录包含岗位名称、公司、薪资、工作城市、发布时间、职位描述等字段。典型问题包括薪资字段“15-25K·14薪”这种文本混排岗位名称里带HTML标签同一公司多条记录写法不统一“经验不限”和“经验不限夹杂全半角空格。拿到数据先别写清洗逻辑跑一遍探查。看形状、字段类型、空值率、重复率再对文本字段做value_counts采样。探查完就知道该写哪些规则了。比如company字段首尾有空格job_title里混着“”标签salary需要拆出最低薪和最高薪“经验不限”这类值是合法的需要保留。import pandas as pd df pd.read_csv(recruitment_raw.csv, encodingutf-8, dtypestr) print(shape:, df.shape) print(columns:, df.columns.tolist()) print(null counts:\n, df.isnull().sum()) print(duplicated rows:, df.duplicated().sum()) print(df[salary].head(20))探查阶段写的代码后面大部分会变成正式清洗规则的雏形。我的习惯是在探查代码里加注释标清楚每段代码在验证哪条规则。比如“验证salary解析后数值是否连续”这样从探查过渡到清洗规则修改时逻辑是连贯的。4.2 第二步把清洗逻辑写成可复用的函数探查完成就开始写正式清洗函数。每个清洗动作封成一个函数比如clean_text处理首尾空格和全半角parse_salary解析薪资区间normalize_city统一城市名。函数化有几个好处规则边界清楚、单测好写、迁移到Spark时逻辑不容易丢。薪资解析是这类数据的核心清洗点。观察数据格式后用正则抽取最低值和最高值再生成平均薪资字段。注意有些记录是“面议”抽取不出来就置为NaN不要强行解析。import re def parse_salary(text): if pd.isna(text): return None, None text str(text).strip() pattern re.compile(r(\d(?:\.\d)?)[Kk]?-(\d(?:\.\d)?)[Kk]?) match pattern.search(text) if not match: return None, None low float(match.group(1)) * 1000 high float(match.group(2)) * 1000 return low, high df[salary_low], df[salary_high] zip(*df[salary].apply(parse_salary)) df[salary_avg] df[[salary_low, salary_high]].mean(axis1)这里有个细节单位是K要乘1000转成元不转的话后续跟其他薪资源对比时直接差一个量级。这就是前面说的一致性问题的具体表现清洗时不把单位统一后面分析出来的平均薪资就是错的。4.3 第三步迁移到Spark跑全量数据小样本规则验证通过后把同样的逻辑迁到Spark。迁移不是把pandas代码原样改语法而是重写为Spark DataFrame的操作。字符串清洗用regexp_replace缺失值用fillna类型转换用cast窗口去重用row_number。下面是Spark版的完整清洗流程示意。from pyspark.sql import SparkSession, functions as F, types as T spark SparkSession.builder.appName(recruitment_clean).enableHiveSupport().getOrCreate() df spark.read.option(header, True).option(inferSchema, False) \ .csv(hdfs:///data/ods/recruitment_raw) # 1. 去重按岗位id和发布时间去重保留最新 from pyspark.sql.window import Window w Window.partitionBy(job_id, publish_date).orderBy(F.desc(etl_time)) df df.withColumn(rn, F.row_number().over(w)).filter(F.col(rn) 1).drop(rn) # 2. 文本字段清理 df df.withColumn(company, F.trim(F.col(company))) \ .withColumn(job_title, F.regexp_replace(F.col(job_title), [^], )) # 3. 薪资解析注册UDF后调用 parse_salary_udf F.udf(lambda s: parse_salary(s), T.StructType([ T.StructField(salary_low, T.DoubleType()), T.StructField(salary_high, T.DoubleType()) ])) df df.withColumn(parsed, parse_salary_udf(F.col(salary))) df df.withColumn(salary_low, F.col(parsed.salary_low)) \ .withColumn(salary_high, F.col(parsed.salary_high)) \ .drop(parsed) # 4. 过滤手机号格式非法的记录 df df.filter(F.col(phone).rlike(^1[3-9][0-9]{9}$)) # 5. 写回DWD层 df.write.mode(overwrite).partitionBy(dt).parquet(hdfs:///data/dwd/recruitment_clean)这段逻辑跟pandas版本一致换的是执行引擎。跑之前先在Spark上对一个月的数据抽样验证对比抽样结果与pandas清洗结果是否一致不一致就说明迁移过程中写错了这一步不能省。集群任务跑全量数据通常要几十分钟甚至几小时等跑完再发现逻辑错误成本太高。4.4 第四步清洗质量校验清洗完成后必须校验不能默认跑完就是对的。我用三个指标做清洗报告记录数变化、主键唯一率、关键字段空值率。清洗前记录100万条清洗后99.2万条多出来的0.8万哪里去了去重删了多少、异常过滤删了多少逐项列清楚。唯一率用主键去重前后的记录数比值衡量空值率看清洗前后是否按预期下降。有些项目会搭建自动化的数据质量检查框架每天定时检查核心表的记录数波动、空值率阈值、主键重复数超阈值就告警。这个框架本身不难核心是一张质量规则配置表和定时调度难的是规则参数的设定。参数设得太松脏数据漏过去设得太紧天天误报最后大家都不看告警了。参数建议先观察两周正常波动范围再取边界值加上缓冲。校验报告最好沉淀成固定格式的文档字段名、清洗前后统计量、异常说明一目了然。做数据治理不是写论文是要让每个下游数据使用者知道这个表能信到什么程度。5. 常见问题与排查技巧实录5.1 高频问题速查表问题见得多了我整理了一张速查表新同学排查时可以先对号入座。现象常见原因排查思路清洗后记录数大幅减少去重主键设得过宽把不同业务记录误判为重复复查subset字段抽样看被删数据数值字段清洗后均值突变异常值删除阈值过激或填充值类型不对对比清洗前后的分位数分布检查fillna的取值时间字段解析后大量为空数据里存在多种日期格式或混入中文日期列出无法解析的样本补充分支解析规则中文乱码原始文件编码与读取编码不一致尝试utf-8、utf-8-sig、gbk逐项确认Spark任务卡顿或OOM按键去重引发数据倾斜或UDF开销过大查看stage耗时尝试用内置函数替代UDF下游指标对不上清洗规则在不同任务中不一致收敛清洗逻辑用统一的数据质量检查框架5.2 几个容易被忽略的坑第一个坑是编码问题。CSV文件从Windows系统导出的经常是GBK编码用utf-8读会直接乱码或者报错。读数据时先确认编码拿不准就多试几种别在清洗阶段就引入新的脏数据。第二个坑是隐式类型转换。pandas里一列既有数字又有字符串时整列会被推断成object类型你以为在做数值运算实际是在拼接字符串。在Snapshot和全量数据处理中都发生过这类问题现在我的习惯是读入时显式指定dtype该是数值的字段绝不偷懒。第三个坑是清洗逻辑分散在脚本里缺乏统一入口。规范做法是把清洗规则收敛到一个模块输入原始表调用各清洗函数输出干净表。这样规则变更、问题排查都只动一个地方调用方只依赖输出不会各自为政地重复清洗。第四个坑是只关注数据不关注元数据。数据字典、字段口径说明、规则版本这些文档比清洗代码本身更能说明问题。我接手过别人留下的清洗脚本没有注释、没有字段说明里面有一个长长的正则没人知道它的含义最后只能靠猜。这种脚本就算跑得对也不敢让它上线跑全量。第五个坑是忽略抽样验证。全量数据清洗跑一遍耗时很久如果不清洗前在小样本上验证规则很容易出现全量任务运行到一半才发现规则写错的情况。特别是正则表达式这种文本规则样本上跑通很容易全量里各种边界样本会不断教育你。我的做法是先抽样一万条清洗看结果再缩放执行。6. 最后几个实操层面的心得数据清洗做了这么多年我最深的体会是清洗工作的核心不是写代码而是判断。判断每条规则该不该写、阈值该定多少、缺失值该删还是该填、异常值该标记还是该修正这些判断依赖对业务逻辑的理解。技术手段都是现成的pandas有fillna和drop_duplicatesSpark有DataFrame APIMapReduce也有一套成熟模板真正决定数据质量的是清洗规则的准确性和可维护性。如果你正在做一个大数据清洗相关的项目我的建议是先把探查这一步做扎实。花半天时间看数据分布省的是后面几天的返工时间。规则一定要写文档哪怕只是几行注释也要让三个月后的自己看得懂。清洗前后各留一份统计报告这既是质量的证明也是和业务方沟通的素材。最后强调一点清洗逻辑尽量集中收敛别散落在各种临时脚本里数据质量是团队资产不是某个人的手工作坊。做清洗这件事耐心比聪明更重要。脏数据永远比你想象的更多模式也永远比你预想的更复杂但只要把探查、规则、执行、校验这套流程走顺大部分问题都能在可控范围内解决。
返回列表