☰
数据清洗实战:从Pandas到Spark的完整方法论
2026/9/30 3:37:14 网站建设 项目流程

干这行这么多年,我越来越觉得数据清洗才是大数据项目里真正见功底的地方。很多人以为搭个集群、跑个分析模型就是大数据,实际上一套数据从采集端到可视化大屏,清洗环节通常要吃掉整个项目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', encoding='utf-8-sig') # 统计空值率,决定删除还是填充 null_ratio = df.isnull().mean() # 数值字段用中位数填充 df['price'] = pd.to_numeric(df['price'], errors='coerce') 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', ascending=False) df = df.drop_duplicates(subset=['order_id'], keep='first')

Spark里同样有dropDuplicates方法,指定subset参数即可。但分布式中要注意一个坑:去重发生在shuffle阶段,如果数据量极大且主键重复率很高,容易引发数据倾斜,个别executor数据量暴增拖垮整个任务。这种情况可以给主键加随机后缀先分散,再去掉后缀做最终去重,不过绝大多数清洗场景用不上这个技巧,知道有这回事就行。

3.3 格式归一:类型、编码、时间、地理坐标

格式归一的目标是让每个字段就一种格式,这是后续计算的基础。我举几个高频场景。

时间字段是最乱的,常见的有“2024/01/05 08:30:00”“2024-01-05 08:30”“2024年1月5日”混在一起。统一用to_datetime解析,解析不了的置为NaT再单独处理。字符串类型的手机号、身份证号,看起来是数字但不能当数值处理,必须在读入时就指定dtype=str,防止前导零丢失。

地理坐标也是个经典坑。GPS设备、地图SDK、第三方接口各自输出不同坐标系,直接按数值计算距离会产生几公里甚至几十公里的误差。清洗时要根据数据来源识别坐标系类型,统一转换后再落库。这个转换逻辑很成熟,网上有现成库,但规则层必须写明每批数据的坐标系来源。

正则表达式是做格式归一最顺手的工具。手机号、邮箱、邮编的校验,都靠正则。跑大数据集群之前,先用正则在小样本上跑一遍,确认没有误伤,再上线全量任务。我见过一个案例,清洗规则把订单备注里的“电话138xxxx8888”全部误判为非法手机号,差点把备注信息清空,就是正则边界没写好。

import re def clean_phone(value): text = str(value).strip() if re.fullmatch(r'1[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'], errors='coerce') # 业务规则:年龄合法范围 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', encoding='utf-8', dtype=str) 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(axis=1)

这里有个细节:单位是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_duplicates,Spark有DataFrame API,MapReduce也有一套成熟模板,真正决定数据质量的是清洗规则的准确性和可维护性。

如果你正在做一个大数据清洗相关的项目,我的建议是先把探查这一步做扎实。花半天时间看数据分布,省的是后面几天的返工时间。规则一定要写文档,哪怕只是几行注释,也要让三个月后的自己看得懂。清洗前后各留一份统计报告,这既是质量的证明,也是和业务方沟通的素材。最后强调一点,清洗逻辑尽量集中收敛,别散落在各种临时脚本里,数据质量是团队资产,不是某个人的手工作坊。

做清洗这件事,耐心比聪明更重要。脏数据永远比你想象的更多,模式也永远比你预想的更复杂,但只要把探查、规则、执行、校验这套流程走顺,大部分问题都能在可控范围内解决。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询