ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

数据清洗实战:从Pandas到Spark的完整流程与关键技术解析

数据清洗实战:从Pandas到Spark的完整流程与关键技术解析 咱们直接进入正题。做了这么多年大数据相关的工作被问得最多的问题不是怎么搭集群、怎么调Spark参数反而是最不起眼的数据清洗怎么搞。这让我挺意外的但仔细想想又在情理之中——很多刚入行的朋友拿着Pandas跑个dropna()就觉得自己会数据清洗了可真到了生产环境面对几千万条脏数据才发现当初的理解有多天真。这篇东西我结合自己带过的一些项目来写包括招聘数据清洗、农产品价格清洗、网约车大数据项目里基于Spark和MapReduce的清洗流程把数据清洗这件事从原理到实操完整捋一遍。不管你是刚接触Pandas的新手还是准备做大数据毕业设计的学生或者已经在搞集群的同学我尽量写得让每个人都能找到自己需要的东西。1. 先把数据清洗这件事说透它到底是什么为什么绕不开1.1 一个被低估的环节数据清洗Data Cleaning说白了就是给数据体检治病。原始数据从业务系统、传感器、日志文件、爬虫程序里出来的时候几乎不可能是干净整齐的——就像刚从菜市场买回来的菜带着泥、烂叶、虫眼你得先择菜洗菜切菜才能下锅。数据清洗就是那个择菜洗菜的过程。很多人觉得清洗只是把空值删掉这么简单这是最大的误解。我在实际项目里见过的脏数据形态五花八门字段错位、编码混乱、重复记录、异常值、时间格式不统一、单位不一致……每一项处理不当后面做统计分析、机器学习建模、可视化展示全都会跑偏。我记得有个做农产品价格分析的项目原始数据里北京新发地市场被写成了北京新发地、新发地市场、bjxfd等多种形式如果不做清洗和标准化按市场维度做聚合统计时这些就会被当成三个不同的市场。类似这样的问题靠直觉是发现不了的。1.2 数据质量的两个层面做数据清洗前先要建立两个层面的认知。第一层是数据质量评估。清洗前你得先摸个底知道数据烂到什么程度。实践中我一般从六个维度检查完整性有没有缺失、准确性值对不对、一致性同一个东西描述是否统一、唯一性有没有重复、时效性数据是否过期、有效性是否符合业务规则。这六项就是数据质量的六边形战士。第二层是清洗策略制定。摸完底之后针对不同类型的问题制定处理方案是删除、填充、转换还是标记都得有明确的规则。比如缺失率超过70%的字段我的原则是直接弃用或者单独拆分因为无论怎么填充这个字段的信息量已经不足以支撑分析而缺失率在5%以内的字段可以考虑中位数、众数或者模型预测填充。1.3 数据清洗在大数据体系中的位置如果站在整个大数据架构的角度看数据清洗是ETLExtract-Transform-Load即抽取-转换-加载流程里最核心的T环节。一个典型的大数据架构通常包括四个层次数据采集层、数据存储层比如HDFS、Hive、数据计算层MapReduce、Spark、数据应用层可视化、报表、机器学习数据清洗卡在采集和存储之间或者存储和计算之间起到的是一个把关人的作用。没有这个把关人上游的数据越脏下游的分析结果就越不可信。行业里有个著名的垃圾进垃圾出Garbage In, Garbage Out原则说的就是这个道理。很多数据分析结果之所以被业务方质疑往往不是算法不够先进而是数据源头就没洗干净。2. 从体检到治理一套完整的数据清洗流程设计2.1 第一步数据探查——动手之前先把脉我见过太多人拿到数据直接开写清洗代码这是非常危险的。正确做法是先做数据探查Data Profiling把数据的底细摸清楚。以Pandas为例我每接到一批新数据固定动作是先跑这几行代码import pandas as pd # 加载数据 df pd.read_csv(raw_data.csv, encodingutf-8) # 查看数据规模 print(df.shape) # 查看字段类型和缺失情况 df.info() # 查看描述性统计 df.describe() # 查看缺失值汇总 missing df.isnull().sum() print(missing[missing 0]) # 查看重复记录 print(df.duplicated().sum()) # 抽样查看具体数据 df.head(10)这一步输出的信息能让你快速判断数据集有多大有没有空值空值集中在哪些字段有没有明显的重复数值型字段的分布是否合理。有些时候光看describe()的输出就能发现异常比如某个价格字段的最小值是负数或者最大值是平均值的几百倍那基本可以断定这里有脏数据。我习惯在探查阶段就做记录把发现的所有问题整理成一个数据质量问题清单一项项列清楚问题字段、问题类型、影响范围、初步判断的处理方案。别嫌这一步麻烦这是整个清洗工作最重要的地基。2.2 第二步问题分类——脏数据的四种典型形态结合大量项目经验我把大数据场景下最常遇到的脏数据归纳为四类第一类是缺失值。这是出现频率最高的问题表现为字段值为空或者为NULL。缺失的原因多种多样可能是业务系统没采集到可能是接口不稳定丢包了也可能是用户压根没填。处理缺失值不能一刀切要分情况讨论个别字段缺失可以用统计值填充关键字段大面积缺失可能需要回源头找原因如果某个字段对分析目标至关重要且缺失严重可能要考虑重新采集。第二类是重复值。同一个实体在数据里出现多次通常是因为数据合并时没有去重或者采集端重复上报。网约车项目里最常见的场景就是订单数据的重复上报同一个订单因为网络问题被推送了两三次如果不按订单ID去重后面算接单率、算平均时长全都会被顶上去。第三类是异常值。异常值不等于错误值它可能是真实发生的极端情况。比如订单时长偶尔出现几百分钟不一定是数据错误可能是长途订单。真正需要处理的异常值是指那些明显超出业务合理范围的值比如负数的价格、超过100%的比例、明显乱码的字符串。判断异常值的方法除了业务经验之外可以用3σ原则或者四分位数IQR法做辅助判断。第四类是格式和一致性问题。这一类最琐碎也最消耗时间。比如日期格式不统一2023/01/01、2023-01-01、20230101、2023年1月1日手机号有带前缀有不带的性别字段有写男Mmale的城市名有带市字有不带的。格式问题不会导致程序报错但会直接影响聚合统计和关联查询的结果。2.3 第三步方案设计——清洗规则要有据可依问题摸清了接下来是制定清洗规则。这一步的核心是每条清洗规则都必须来自业务逻辑不能拍脑袋。我举一个实际经历过的例子。在招聘数据清洗项目里薪资字段经常是8k-15k这样的文本要做分析就必须拆分成最低薪资和最高薪资两个数值字段。拆分规则的确定就需要业务知识这里的k代表千·14薪代表一年发14个月工资。没有这个背景知识光靠正则表达式是拆不清楚的。清洗方案设计还有一个要点是保留溯源。我不建议直接覆盖原始数据而是把清洗过程做成标准化流程每一步都记录操作日志。比如新建一列clean_status标记这条记录是否被清洗过、清洗规则是什么。这样即使下游发现问题也能回查清洗逻辑而不是面对一堆已经无法还原的数据干瞪眼。2.4 第四步执行与验证——清洗完要复查清洗流程跑完之后必须做验证。我见过不少同学清洗完数据就万事大吉结果汇总统计出来的数字跟业务方手里的数对不上来回返工。验证我一般做三层第一层是基础校验看清洗前后的数据量变化、缺失率变化、去重数量是否符合预期第二层是抽样人工检查随机抽几百条记录人眼比对一下清洗结果是否合理第三层是下游验证拿清洗后的数据跑一遍关键指标和业务方确认数值是否合理。三层验证都通过这批数据才算出师。3. 核心细节解析手把手拆解数据清洗的关键技术点3.1 缺失值处理的三板斧缺失值处理主流方法就三类删除、填充、保留。删除是最简单粗暴的方式适用于缺失比例极低我通常以5%为界且删除后不影响样本代表性的场景。调用方式也简单# 删除带有缺失值的行 df_cleaned df.dropna() # 只删除指定列存在缺失的行 df_cleaned df.dropna(subset[order_id, price])但删除有一个大坑如果你删掉的缺失记录恰好集中在某个特定群体比如某个城市的订单数据几乎全是空的那删完之后样本就偏了分析结果会以偏概全。所以删除之前一定要先分组看看缺失分布。填充是更常用的手段。数值型字段我倾向于用中位数填充而不是均值因为均值容易受极端值影响中位数更稳健。分类型字段用众数填充也就是出现最多的那个类别。时间序列数据可以用前向填充ffill或者插值法# 中位数填充 df[price].fillna(df[price].median(), inplaceTrue) # 众数填充 df[city].fillna(df[city].mode()[0], inplaceTrue) # 前向填充适用于时间序列 df[value].fillna(methodffill, inplaceTrue) # 线性插值 df[value].interpolate(methodlinear, inplaceTrue)保留听起来奇怪但实际很常见。有些业务场景里缺失值本身就是一种信号。比如用户调研里未填写收入可能暗示用户对收入敏感这时候与其强行填充一个尴尬的数值不如把缺失本身作为一档单独处理在建模时增加一个是否缺失的指示变量。3.2 重复值处理去重不是那么简单Pandas里去重一行代码就行df_cleaned df.drop_duplicates(subset[order_id])但这背后有几个容易踩的坑。第一个坑是去重键的选择。如果只按部分字段去重要小心误删。比如网约车订单可能有重名的情况如果订单号因为系统bug重复生成但实际是两个不同订单按订单号去重就会误删。所以去重前最好结合业务知识判断哪些字段组合能唯一标识一条记录。第二个坑是保留哪条记录。多胞胎记录里可能有的完整有的残缺有的更新有的过期。Pandas的drop_duplicates默认保留第一条但实践中我经常需要保留最新的一条或者最完整的一条。这时可以先用sorted_values排序再按关键字段去重就能做到保最新# 先按时间倒序排列再按订单号去重保留时间最新的 df_sorted df.sort_values(create_time, ascendingFalse) df_cleaned df_sorted.drop_duplicates(subset[order_id], keepfirst)第三个坑是相似重复。有些记录不是完全一样而是高度相似比如公司名阿里巴巴和阿里巴巴集团这在ID-Mapping时才需要处理但如果你直接按公司名字符串去重两条都会被保留。判断相似重复需要用到编辑距离或者SimHash这类文本相似度算法这里不展开但你要知道这个问题的存在。3.3 异常值检测的三种实用方法异常值检测在有监督和无监督学习里都有大量算法但实际做数据清洗我用得最多的是这三种第一种是3σ原则。假设数据服从正态分布超出均值±3倍标准差的点视为异常。代码实现import numpy as np def detect_outlier_3sigma(data): mean data.mean() std data.std() lower_bound mean - 3 * std upper_bound mean 3 * std return (data lower_bound) | (data upper_bound) df[is_outlier] detect_outlier_3sigma(df[price])第二种是IQR四分位距法。用Q125%分位数和Q375%分位数定义正常范围[Q1 - 1.5 * IQR, Q3 1.5 * IQR]超出即为异常。IQR法对偏态分布更稳健是我个人最常用的一种。第三种是业务规则判断也是最容易被忽视的。很多异常从统计角度看是离群值但业务上完全合理反过来有些统计上完全正常的值业务上一看就是错的。比如年龄字段200岁统计上它是个离群点没错但年龄30这个值统计上正常如果某条记录生日填错了导致计算出年龄150也需要被识别。所以我的建议永远是统计方法做初筛业务规则做终审。异常值找到之后的处理策略要回到业务目标。做描述性统计异常值可以单独分析做回归模型异常值可能严重影响拟合效果做聚类异常值可能形成独立的簇。不能一概而论地删除我通常先标记出来单独观察再决策。3.4 格式统一与标准化最琐碎但也最影响分析结果格式统一这块主要以时间格式、字符串清洗、单位换算三大类为主。时间格式统一推荐统一用Pandas的to_datetime# 多种格式混合的时间列统一成标准格式 df[time] pd.to_datetime(df[time], formatmixed, errorscoerce)这里的errorscoerce很关键解析失败的值会变成NaT方便后续单独排查。注意上面参数formatmixed是从Pandas 2.0开始支持的写法旧版本可以考虑用infer_datetime_formatTrue。字符串清洗的内容更多去空格、去换行、统一大小写、替换不规范写法。比如手机号字段可能存在139 1234 5678这种带空格的情况需要去掉再统一为11位数字。英文姓名有的大小写不统一需要统一为规范格式。这些操作可以用Pandas的str方法链式处理# 去除首尾和中间空格统一小写 df[name] df[name].str.strip().str.replace( , ).str.lower()单位换算最经典的坑是万和元的混用。招聘薪资里有15k、1.5万、15000不统一换算成元/月的话下游计算平均薪资就是灾难。我用正则先提取数字再判断单位完成换算import re def parse_salary(text): if pd.isna(text): return None match re.search(r(\d\.?\d*)\s*([kK万wW]?), str(text)) if not match: return None num float(match.group(1)) unit match.group(2) if unit.lower() k: return num * 1000 elif unit in [万, w, W]: return num * 10000 else: return num4. 从小数据到大数据的跨越Pandas、Spark和MapReduce怎么选4.1 Pandas单机王者但别让它干重活Pandas做数据清洗确实方便语法直观、生态丰富、调试方便我个人在探索期、小数据量的场景下都是首选。但要注意的是Pandas的数据全部加载进内存当数据量达到几百GB甚至TB级别的时候单机内存早就撑不住了。我建议的分界线是数据量在几个GB以内单机内存够用用Pandas数据量大到单机跑不动或者需要分布式处理果断上Spark。Pandas还有一个好搭档是pandas的并行加速库比如modin或者polars它们在API层面尽量兼容Pandas底层做了并行化处理在小集群环境里可以替代Pandas处理更大数据量。但如果说要处理的是真正的大数据还得看Spark。4.2 Spark分布式清洗的正确姿势在网约车大数据项目里我做过基于Spark的完整数据清洗。Spark的核心思想是把数据切分成多个分区partition分发到集群的不同节点上并行处理自家兄弟Spark SQL用起来尤其方便。用Spark做清洗我通常会写类似这样的代码from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, isnan, isnull, count # 初始化SparkSession spark SparkSession.builder \ .appName(DataCleaning) \ .getOrCreate() # 读取数据 df spark.read.csv(hdfs:///data/ride_orders.csv, headerTrue, inferSchemaTrue) # 查看缺失情况 df.select([count(when(isnull(c), c)).alias(c) for c in df.columns]).show() # 填充缺失值 df_cleaned df.fillna({price: 0, city: 未知}) # 去重 df_cleaned df_cleaned.dropDuplicates([order_id])注意Spark的处理方式和Pandas有很大区别Spark是惰性求值代码写完后只有真正触发show()、count()、write这些动作时才真正执行计算。这个特性让Spark可以做优化调度但初学者容易犯的错就是在每个步骤后面都加一个show()导致整个任务反复跑集群资源全浪费在重复计算上了。另一个Spark清洗的关键点是分区数控制。分区太多会导致调度开销大太少会导致数据倾斜。我一般按照每个分区处理128MB左右数据的经验值来设置分区数数据量不确定时可以先用df.repartition()做调整再检查Spark UI里每阶段的任务执行时间观察是否有明显的长尾任务。基于MapReduce的清洗方案在早期的Hadoop生态里很流行用Java写Mapper和Reducer做数据清洗。它的优势在于稳定、可控、不依赖内存缺点是开发效率低。现在我更推荐直接用Hive SQL或者Spark SQL同样的逻辑用SQL写几行就完成了。如果毕业设计或者面试需要展示MapReduce功底写一个清洗任务作为案例没问题但真实生产环境里我个人已经很少用纯MapReduce做清洗了。4.3 大数据集群部署与清洗任务的调度衔接既然提到了大数据清洗就不能不谈它的运行环境。网约车这个项目里我们是三台节点搭建的Hadoop集群NameNode和ResourceManager部署在一台主力节点上另外两台做DataNode和NodeManager。Spark跑在YARN上资源分配根据数据量动态调整。集群部署看起来只是把组件装起来但里面坑不少JAVA_HOME配置、节点之间免密登录、各组件版本兼容性每一样都能折腾人半天。清洗任务一般放在凌晨跑通过调度工具每天定时执行。这一步的价值在于数据清洗不是一次性工作而是流水线上每天都要重复的日常动作。把清洗流程固化成调度任务数据到点自动清洗入库第二天业务拿到的就是干净数据。这一套跑顺之后最直观的感受就是下游分析不再天天救火了。5. 实战项目拆解三个典型场景的数据清洗全流程5.1 招聘数据清洗处理最典型的文本型脏数据以2020年MathorCup大数据挑战赛赛道A这类招聘数据分析项目为例原始数据包含公司名称、岗位名称、薪资范围、学历要求、工作经验、城市、发布时间等字段。这类数据的典型问题我之前提过薪资格式混乱除此之外还有几个高频问题。公司规模字段经常是1000-9999人、10000人以上这样区间文本要转换成有序数值特征可以把下限作为数值[1000, 9999]取100010000人以上取10000。学历要求则要统一映射本科、本科及以上、本科及同等学力都归到本科这套映射关系要单独维护一张字典表。完整的清洗流程做成代码的话我提供一个简化版的Pipeline框架def clean_recruitment_data(raw_df): # 1. 去除全空行 df raw_df.dropna(howall) # 2. 去除关键字段为空的行 df df.dropna(subset[job_name, company_name]) # 3. 去重 df df.drop_duplicates(subset[job_id]) # 4. 解析薪资 df[salary_min] df[salary].apply(lambda x: parse_salary(x)[0]) df[salary_max] df[salary].apply(lambda x: parse_salary(x)[1]) # 5. 标准化学历 df[education] df[education].map(edu_mapping) # 6. 过滤异常值 df df[(df[salary_min] 0) (df[salary_max] df[salary_min])] return df这里有个容易忽略的细节过滤salary_max salary_min这一步很多人不做导致下游计算平均薪资时出现负值区间分析结果完全不可用。像这种逻辑异常值光靠统计方法是检测不出来的必须结合业务规则。5.2 农产品价格数据清洗序列化的时间维度陷阱农产品价格数据最典型的特点是时间序列结构。每天、每个市场、每个品种有一行价格记录清洗的重点除了常见的缺失和异常还有一个独特问题节假日和休市日。很多农贸批发市场周末可能不更新价格数据里对应的日期就没有记录如果用简单的fillna(methodffill)前向填充会把周四的价格平移到周六看似合理但实际上掩盖了周末休市的事实。更合理的做法是维护一个交易日历只在交易日之间做插值或者前向填充非交易日保持为空或者标记为休市。另一个常见问题是价格波动异常。蔬菜水果价格受天气影响大单日暴涨暴跌不一定错。但如果同一天、同一市场、同一品种出现两条价格相差10倍以上的记录大概率是录入时小数点打错了位置。我的处理方式是加一条规则同一(market, product, date)组合内如果有记录偏离该品种过去7天中位数的3倍以上标记出来人工复核。用Pandas处理这类数据时还有个性能细节如果数据是逐日的日期字段建议直接设置为索引这样resample重采样和滑窗计算都方便。代码示例df[date] pd.to_datetime(df[date]) df df.set_index(date) # 按品种和市场分组对价格做7天滚动中位数 df[price_med_7d] df.groupby([market, product])[price] \ .transform(lambda x: x.rolling(window7, min_periods3).median())注意groupby和rolling的组合方式transform保证输出和原数据行数一致方便后面做异常比对。5.3 网约车大数据项目基于Spark的亿级订单清洗最后聊聊网约车项目这是我最完整的分布式清洗实践。数据量到了亿级单机Pandas已经无能为力必须上Spark。整体清洗任务的Pipeline我大致分成四段第一段数据接入与格式解析。原始数据多半是JSON格式的日志或者CSV文件。用Spark直接读取后第一步是解析嵌套字段比如把location这种JSON字符串拆成lng和lat两列。这一步用from_json配合schema定义最稳妥解析失败的行记录到错误日志表。第二段缺失与去重。订单ID是核心键缺失的直接剔除重复的按照create_time排序保留最新。司机ID和乘客ID的缺失则选择保留但标记因为这两列可能用于后续关联维表缺失的行可以作为匿名订单单独分析。第三段业务异常值清洗。定义规则订单金额必须为正且不超过某个上限比如5000元超过可能是专车长单也可能是录入错误订单时长在1分钟到24小时之间经纬度必须在中国范围内。违反规则的记录不是直接删除而是写入异常明细表将来如果业务方想排查风控问题这些数据还能派上用场。第四段字段标准化与特征衍生。把时间戳统一转换为日期、小时等维度字段计算订单时长、司机收入分成比等衍生字段。这一步是为了下游分析直接使用避免每个分析任务重复写转换逻辑。这里要特别强调清洗任务的幂等性。因为清洗是每天定时跑的如果某天任务中途失败重跑或者上游数据被重新推送清洗结果不能出现重复或遗漏。我的做法是使用overwrite模式写入结果表并且清洗任务入口先做删除当天分区再写入。这样即使重复执行结果也是一致的。6. 清洗路上的拦路虎常见问题与排查技巧实录6.1 数据量一大Pandas内存溢出怎么办这是被问得最多的问题。解决方案按优先级排列先看能不能只读需要的列usecols参数减少内存占用再看能不能把数据类型降级比如把int64转成int32甚至int8把object类型转成category类型还不行就换chunksize分块处理终极方案是上Spark。有个我用了很久的小技巧用Pandas读取大数据文件时先不指定dtype用df.info(memory_usagedeep)查看内存占用找出占用最大的列再针对性地做类型优化。很多CSV里的数值列因为混入了空值被Pandas读成了float64其实完全可以在读取时用dtype强制指定内存立刻节省一半。6.2 清洗前后数据量对不上怎么排查出现这种情况九成是去重或者过滤条件设置得太狠了。排查思路是先分别统计缺失删除、去重删除、异常删除三个环节各自影响了多少行把每个环节的删除量记录下来。再抽样对比被删除的几条记录看看是合理删除还是误删。最后如果删除比例超过10%我会停下来反思规则是不是太激进。6.3 Spark任务跑得慢怎么定位瓶颈Spark任务慢的常见原因数据倾斜、分区过大、Shuffle过多。排查时先看Spark UI找耗时最长的Stage看它的Shuffle Read大小和Task耗时分布。如果是明显的数据倾斜某个Task耗时是其他Task的几十倍优先考虑加盐salting或者换用repartition均衡分区如果是Shuffle过大看能不能通过broadcast join替代sort merge join减少网络传输。6.4 清洗规则写完结果还是有问题我复盘过很多次这种情况最后发现普遍原因是规则之间的执行顺序出了问题。比如先做异常值过滤再做缺失值填充和先填充再过滤结果完全不同。我总结的经验是先做格式统一和标准化再做缺失处理再做去重最后做异常值过滤。这个顺序不是拍脑袋定的而是因为格式问题和缺失问题会影响去重时是否重复的判断而异常值判断又依赖标准化后的干净数据。数据清洗这个领域看起来门槛低好像会几个函数就算会了但真正做深了会发现它考验的是数据敏感度、业务理解力和工程落地能力的综合水平。每次接到一批新数据我仍然会保持着先探查、后设计、再执行、最后验证的习惯。按部就班看起来慢实际上才是最快、最稳的路。最后分享一个小经验清洗规则的代码一定要写注释、写文档把每一条规则的业务依据写清楚。做数据这行人换了一拨又一拨数据还在那里。你留下的清洗代码和规则文档就是后来人最宝贵的工具。这比我上面写的任何一条技术细节都重要。
返回列表