ARTICLE DETAIL

资讯详情

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

数据清洗实战:从pandas到DataX的完整指南

数据清洗实战:从pandas到DataX的完整指南 接到手里的数据十有八九是没法直接用的这不是夸张是常态。做数据分析、跑模型、出报表真正花时间的地方往往不在算法和可视化而是耗在“把数据收拾到能看”这一步上。数据清洗Data Cleansing干的就是这个活把缺失的补上、错误的纠正、重复的去掉、格式不统一的收拾齐整让数据从“能打开”变成“能信、能用、能算”。这篇东西我不打算写成教科书式的概念科普而是按我这些年实际摸爬滚打的经验来聊。内容会覆盖三个层面一是清洗的整体思路和通用流程让你拿到任何一份脏数据都有章可循二是基于 pandas 的手工清洗实操适合数据量不大、需要精细处理的场景三是 DataX 这类批处理工具在数据同步链路里做清洗的玩法适合生产环境里天天跑的活儿。最后单独讲一下工业传感器数据清洗——这是我个人认为“最脏”的数据类型没有之一。无论你是刚入行的分析师、数据工程师还是做物联网平台开发的这篇文章里总有一部分能用得上。1. 数据清洗先搞清楚到底在清什么很多人一上来就写代码这是最大的误区。你连数据是怎么脏的都没弄明白写出来的清洗逻辑大概率是拍脑袋拍完还得返工。我习惯先做一次“数据勘察”把脏数据的类型摸清楚再决定用什么手段处理。1.1 脏数据从哪来不是所有“脏”都长一样我经常把脏数据分成六类每一类的成因和处理思路都完全不同缺失值这最常见单元格是空的或者填了“N/A”“NULL”“-”这类占位符。成因一般是采集漏了、接口返回字段为空、系统迁移时丢了。处理思路是删除、填充、还是单独标记得看业务含义不能一刀切。重复数据同一件事被记了两遍甚至更多遍。成因很杂比如数据重复上报、多表关联时产生了笛卡尔积、ETL任务重复跑了一次但没做幂等。需要注意的是完全相同的重复好去就怕“看起来不一样其实是同一件事”的重复例如“张三”和“张 三”、“北京”和“北京市”。格式不一致同一个字段里幺蛾子最多。日期有“20230101”“2023-01-01”“2023/1/1”各种写法手机号有的带86前缀、有的带空格金额有“1000”“1,000”“壹仟元”。这类问题必须在进入分析前统一否则一排序、一聚合就翻车。异常值超出合理范围的值。比如人的年龄填了230传感器温度读出来是-9999订单金额是负数。异常值不一定是错也可能代表着真实的异常事件比如设备故障、用户退单所以要分清楚是“录入错误”还是“真实离群”。逻辑矛盾最隐蔽的一类。比如注册时间晚于最后登录时间、订单状态标了“已发货”但没有物流单号、身份证号码里的出生日期和生日字段对不上。机器检查很难全发现多半得靠业务规则。编码与字符问题中文乱码、全角半角混用、\n和回车符混在文本里、BOM头干扰列名识别。这类问题不处理轻则匹配不上重则整个文件读不进来。还有一个经常被忽略的来源多表关联时字段语义对不上。小A表里的“客户ID”在小B表里叫“user_id”一边是字符串一边是整数一边是内部编号一边是业务编号——这种清洗不是修单元格而是要重建字段映射关系属于广义的数据一致性清洗。1.2 清洗不是“一把梭”而是五段式流程我个人习惯把清洗拆成五个阶段每一步都有明确产出这样既不会漏也方便跟同事对齐进度数据探查拿到数据先不动手用info()、describe()、nunique()这类函数把数据集的形状、字段类型、缺失情况、取值分布摸一遍。目标是在动手前就意识到表多大、哪些列隐患大、哪些列基本是废列可以直接丢。规则制定根据业务含义为每个重点字段制定清洗规则。比如“年龄字段超过100的按缺失处理”“金额字段小于0的先查退款记录再定”“时间字段统一解析为YYYY-MM-DD HH:MM:SS”。规则一定要写成文档哪怕是临时笔记这个习惯在项目交接时能救命。执行清洗按规则写代码处理先复制一份原始数据再动手原始数据始终保持只读。这一步里我会把处理逻辑封装成函数处理的每一步都打印变更前后的行数保证每步都可追溯。质量验证清洗完不等于完事。用交叉统计验证清洗后唯一值数量对不对、各组分布比例是否异常、抽样50条肉眼复核一遍。我还会跑一遍下游常用逻辑比如按月汇总、按用户去重看看输出结果是否符合常识。回归与留痕把清洗脚本连同参数、执行时间、影响行数一起记录下来。数据是活的明天新数据进来同一个脚本还得再跑不留痕等于白干。这五步看起来简单但大多数清洗项目翻车都翻在第一步和第二步——没探查清楚就动手规则又没和业务对齐清洗完数据比原来还难用。2. pandas清洗实操从DataFrame到干净数据pandas是Python生态里做结构化数据清洗最顺手的工具没有之一。这一节我用一个客户订单表为例按实际操作的顺序讲一遍核心代码和背后的判断逻辑。示例数据包含customer_id、name、phone、province、order_date、amount、status这些字段。2.1 读进来先别急着跑dtype、索引与缺失值普查很多同学习惯pd.read_csv(data.csv)接一行df.head()就开干省掉的几步恰恰是后面翻车的根源。import pandas as pd df pd.read_csv( orders.csv, encodingutf-8-sig, # 处理BOM避免第一列列名乱码 dtype{customer_id: str}, # ID按字符串读防止前导0丢失 parse_dates[order_date], # 日期列直接解析 )读进来之后这三行必看df.info() # 每列的非空数量、dtype扫一眼就知道哪些列是重灾区 df.describe() # 数值列的分布min/max异常一眼可见 print(df.isnull().sum()) # 精确的缺失统计这里有一个新手容易踩的坑customer_id不要用默认方式读。Excel里的00123经过pandas读取后可能变成123数字ID如果后面要做字符匹配直接对不上。所以读的时候指定dtype{customer_id: str}是常态操作。同理手机号、身份证号这类不该参与数值计算的字段都要显式按字符串读。2.2 缺失值处理的三种思路按业务场景选缺失值不是只能删或者只能填完整的选择是“删除、填充、标记”三选一具体怎么选要看这列数据是干什么用的。**删除适用于缺失比例过高或无分析价值的列。**如果一列有超过70%都是空的除非它是有特定含义的稀疏标记比如“退单原因”否则它对分析结果只有噪音没有信号建议直接删。删除行的场景主要是关键字段为空且后端表也补不回来比如订单表里连订单ID都丢了这种记录留着没有任何用处。**填充适用于数值型字段且有合理填补策略。**最常用的是中位数、均值、众数、前后值填充但很多人忽略了一个前提你用什么值填取决于数据分布和业务逻辑。收入字段严重右偏时用均值填充会被极少数高收入人群拉高中位数更稳时间序列里的缺失用前向填充ffill通常比全局均值更合理因为相邻时刻的值往往更接近。# 金额缺失用该省份客户的金额中位数填充 df[amount] df[amount].fillna( df.groupby(province)[amount].transform(median) ) # 时间序列演示前向填充订单日期缺失的用上一个客户的日期兜底 df[order_date] df[order_date].ffill()**标记适用于“缺失本身代表一种状态”的场景。**比如营销表里的“点击时间”为空不代表数据丢失而是代表用户压根没点过。这时候把缺失值填成0或某个特定值反而扭曲了语义。正确的做法是保留空值或者加一列is_clicked作为分类特征。2.3 重复值、格式统一与类型转换的实战细节重复值处理很多人就drop_duplicates()一把梭。问题在于完全一样的记录你当然可以直接删但业务上更常见的是“某些关键字段重复其他字段不同”。比如一个用户下了两单订单字段都不同用户字段相同——这算不算重复当然不算。所以一定要用subset参数指定判断重复的字段集合# 对同一天、同一个客户、同一个金额的订单去重保留最新一条 df df.drop_duplicates( subset[customer_id, order_date, amount], keeplast, )格式统一是脏数据重灾区尤其手机号、文本列。我写过一套“统一手机号格式”的流程先去空格和全角字符去掉86/86前缀统一为11位最后用正则校验。这一步看着麻烦但省掉了后面接各种短信平台、CRM系统时的一堆幺蛾子。df[phone] ( df[phone] .astype(str) .str.replace(r\D, , regexTrue) # 去掉所有非数字字符 .str.replace(r^86, , regexTrue) # 去掉86前缀 ) # 校验筛选出长度不是11位的大概率是脏数据 bad_phone df[df[phone].str.len() ! 11] print(f异常手机号数量: {len(bad_phone)})类型转换的坑更多。字符串转数值时1,000直接astype(float)会报错得先去掉千分位逗号bool列转int时True/False在pandas 2.x下的表现跟老版本有差异日期字符串格式不统一时to_datetime会飘红这种情况要么用format参数指定要么先做一次格式规整再转。# 处理金额列中的千分位逗号和人民币符号 df[amount_clean] ( df[amount] .astype(str) .str.replace(¥, , regexFalse) .str.replace(,, , regexFalse) .astype(float) )3. DataX批处理清洗让清洗进入生产流水线pandas处理几万、几十万行的数据很舒服但到了每天几千万行的数据同步场景单机跑DataFrame就力不从心了。这时候我需要的是能挂在调度平台上的批处理工具DataX就是我用得比较顺手的一个。它本身是异构数据源同步工具核心能力是把数据从MySQL、Oracle、HDFS、Hive等地方搬来搬去但它的Transformer机制可以在搬运过程中顺便做清洗。3.1 为什么单机脚本不够用先说痛点。在数据仓库项目里最常规的需求是每天凌晨把业务库的增量数据同步到数仓ODS层。用pandas脚本处理的话有四个让人头大的问题数据量大单表几千万行pandas全量load进内存服务器内存分分钟爆掉。异构数据源源库是Oracle目标库是Hive中间还有SQLServer的老系统靠写连接器逐套对接维护成本高到崩溃。并发与调度生产环境里几十张表要同时同步每张表一个pandas进程光进程管理就够受的。一旦某张表失败重跑逻辑还得自己写。资源利用pandas的处理是单机的机器配置再高也只能垂直扩展撑不住横向扩展的流量。DataX解决的是同步的骨架问题它用插件机制对接各种数据源由调度平台统一管理任务跑在分布式环境里天然支持水平扩展。而数据清洗这件事正好可以借用它的Transformer机制在同步过程中顺手完成。3.2 DataX核心机制与本地跑通DataX的工作模型很清晰一个作业Job分成若干个Task每个Task由三个核心组件构成——Reader负责从源端读取、Transformer负责行级数据处理、Writer负责写入目标端。你配置一个JSON文件描述“从哪读、做什么处理、写到哪”DataX的框架就帮你把并发、重试、断点这些都管好了。一个最简单的DataX任务配置长这样{ job: { content: [ { reader: { name: mysqlreader, parameter: { username: etl_user, password: ******, column: [customer_id, phone, order_date, amount, status], connection: [ { jdbcUrl: [jdbc:mysql://192.168.1.100:3306/business], table: [orders] } ] } }, writer: { name: hdfswriter, parameter: { defaultFS: hdfs://nameservice1, fileType: text, path: /warehouse/ods/orders, writeMode: append, fieldDelimiter: \t, column: [ {name: customer_id, type: string}, {name: phone, type: string}, {name: order_date, type: string}, {name: amount, type: double}, {name: status, type: string} ] } } } ], setting: { speed: { channel: 4 } } } }配置本身不复杂真正有讲究的是channel数。它决定了并发度调太低了同步慢调太高了会给源库造成查询压力。我一般先按源表行数预估千万级以内的表先跑4个channel试试观察源库的CPU和IO负载再往上加不建议一上来就开16甚至32。3.3 用Transformer在同步中完成轻度清洗DataX内置了几种Transformer包括dx_replace正则替换、dx_substr截取、dx_filter行过滤、dx_groovy写Groovy脚本做任意处理。光靠这几个就能覆盖掉一部分常规清洗需求。比如订单状态字段里混了已支付 带空格、已支付。带了中文句号直接在写Hive前统一掉transformer: [ { name: dx_replace, parameter: { columnName: status, replaceWith: 已支付, judgeValue: [已支付 , 已支付。, 已支付] } }, { name: dx_filter, parameter: { condition: amount 0 } } ]dx_filter那一段表示只同步金额大于0的订单负数金额的直接在管道里过滤掉下游就不用来回排查了。但说实话内建Transformer的表达式能力有限碰到“需要关联另一张表来补字段”这类清洗光靠它不够。我的做法是DataX负责同步和粗清洗复杂清洗逻辑放到下游的SQL或Spark任务里做。各层各司其职不要试图把所有事情都塞进一个工具。有一个原则我踩坑后总结出来的同步管道里的Transformer越简单越好复杂到需要调试半天的逻辑就别放在同步链路里否则每次同步任务失败了你都不知道是源库的问题、网络的问题还是Transformer表达式写错了。4. 工业传感器数据清洗最脏场景的实战经验做过互联网业务数据再去做工业传感器数据清洗那感觉就像从客厅进了煤窑。业务数据再脏顶多是缺字段、格式乱传感器的数据是真真切切的物理世界噪声信号漂移、设备停机、通信断断续续、量纲五花八门。这一节单独拿出来讲是因为它跟前面讲的表格式数据清洗在方法论上有本质差异。4.1 传感器数据为什么特殊时间序列的脏法不一样传感器数据的第一个核心特征它有天然的顺序属性。一个温度传感器每秒钟上报一次数据今天早上10:00:01来了10:00:02可能就断了10:00:03又回来了。这个“断”在业务表里可能表现为缺失行但当你做趋势分析时缺的不是一行而是一个时间段——清洗时必须考虑时间窗内的连续性而不能像处理客户表那样简单地删掉一行。第二个特征是脏值往往有模式。传感器数据里的异常值不是随机冒出来的常见的有这么几类数值饱和传感器量程上限是100℃读出来一直是999或-9999那是超量程标记。阶跃跳变正常温度是20℃上下波动某条记录突然跳到80℃下一跳又回到21℃。这大概率是传感器受到瞬时干扰或者前端信号处理出了毛刺。斜率异常温度在10秒内上升了50℃物理上不可能即便读出来的每个点本身都在合理量程内。滞后/冻结连续N条记录数值完全一样可能是传感器卡死或设备停机了。时间戳紊乱设备断电重启后内部时钟没同步上报的数据时间戳忽前忽后甚至出现2030年的“未来时间”。4.2 时间戳对齐与异常值判定的工程经验处理传感器数据的第一步永远是时间戳对齐。设备上报频率可能不固定启动时快到毫秒级稳定后慢到秒级分析前必须先重采样到统一时间基准。import pandas as pd # 原始数据ts列为时间戳val列为数值 df[ts] pd.to_datetime(df[ts], utcTrue, errorscoerce) df df.dropna(subset[ts]) df df.set_index(ts).sort_index() # 重采样到1秒间隔缺失值生成NaN后再插值 df_aligned df[val].resample(1S).mean() df_aligned df_aligned.interpolate(methodtime, limit_directionboth)resample里我选了mean聚合是为了处理一秒内有多条记录的情况。然后interpolate(methodtime)是按时间间隔做线性插值比简单的ffill更能保留趋势。但要注意插值只适用于短时间缺失如果设备停了整整一小时线性插值会把这段补成一条斜线对后续诊断完全没意义。所以正确做法是先设定阈值缺失时长超过比如5分钟这一整段就标记为“设备停机”不打补丁。异常值判定方面除了设定物理上下限温度不可能零下100℃我更推荐用滑动窗口做局部异常检测。全局限值能抓出明显错误但抓不住“局部跳变”。最简单有效的方法是# 计算每个点的局部均值和标准差窗口30秒 rolling_mean df[val].rolling(window30, centerTrue).mean() rolling_std df[val].rolling(window30, centerTrue).std() # 超过局部均值±3倍标准差记为异常 df[is_anomaly] (df[val] - rolling_mean).abs() 3 * rolling_std3倍标准差是个经验值具体怎么调取决于业务容忍度。报警类应用宁可多报2.5倍趋势分析宁可少报4倍。我建议先用3倍跑一遍统计异常比例如果异常占比超过3%基本可以断定窗口或阈值设置不合理而不是现场真的有这么多故障。4.3 状态标记比删除更重要传感器数据的清洗观做业务数据清洗时我们的目标是让数据“变干净”。但传感器数据清洗我的核心原则是不要轻易删除任何一条记录而是要给它打状态标签。原因有三第一传感器的异常值可能本身就是最重要的事件信号。一次振动传感器记录的振幅尖峰在清洗阶段被当成“异常”删掉了后面的故障诊断就无法复现当时工况。第二删除会破坏时间序列的连续性下游做频谱分析、特征提取时会踩坑——缺失片段和标记异常片段对算法的含义完全不同。第三清洗过程要可审计。设备供应商半夜甩锅说“你的清洗逻辑把我正常数据都删了”如果你保留原始数据清洗标记可以直接拉出标记字段对质。所以我在工程上给传感器数据加的字段通常是is_valid是否通过质量检查0/1quality_code标记问题类型1正常2超量程3跳变4冻结5设备停机signal_lost数据缺失时间段True/Falsets_repaired时间戳是否经过校正下游使用时is_valid1的数据可以直接进分析模型is_valid0的数据单独存放留给设备组做故障定位。数据库的存储量可能多出10%但换来的是整个过程的可追溯非常值。5. 常见问题与排查技巧实录数据清洗这么多年踩过的坑比写过的代码还多。这一节我把最典型的几个场景整理出来做成一份“问题速查表”你在实操中碰到类似的可以直接照着排查。5.1 六个高频翻车现场速查现象常见原因排查思路read_csv读出来第一列列名多一个\ufeff文件带UTF-8 BOM头改用encodingutf-8-sig读取日期解析大面积报错多种日期格式混用或有非日期文本用errorscoerce定位坏值先看格式分布再统一解析规则drop_duplicates后行数比预期少很多判断重复的字段选窄了不同业务实体被误判为重复重新确认业务唯一键加上keep参数测试不同保留策略聚合结果比手工算的差很多有字符串形式的数值列如1000不是1000没转类型df.dtypes查一遍astype(float)前先清洗千分位和货币符号清洗后下游报表数据翻倍多表关联产生了重复行清洗阶段没有在关联前先去重关联前先按业务键去重关联后再duplicated()检查一次传感器数据插值后出现“悬崖”对长时间缺失段做了线性插值只有短时间缺失才插值长时间缺失整段标记为停机这里面的“清洗后数据翻倍”是我见过最坑的。有一次同事跑数仓任务前一天报表还正常第二天突然所有订单金额都翻了一倍。排查了半天发现他新加一个关联表源表里有历史重复记录关联之后订单行数直接翻倍。后来我们把它固化成了规范任何多表关联的操作之前必须先确认两个表的粒度关联完后立刻检查行数变化是否在合理范围内。5.2 数据清洗的纪律性清单工具和方法论都讲了最后分享一份我自己的干活清单。“纪律不是束缚是在你熬夜调数据时帮你兜底的保险丝”永远保留原始数据。不管原始数据多脏全量备份一份只读副本清洗脚本处理的是副本。宁可多花一倍存储不要因为误删后悔。每一步清洗都要打印变更记录。处理前多少行、处理后多少行、删了多少重复、填了多少缺失这些数字是验证逻辑是否正确的最快途径。数据清洗如果跑完输出“成功”就完事那等于没做。规则先成文再写代码。哪怕规则只写在笔记软件里也好过直接写代码。因为清洗规则本质上是业务逻辑代码是给机器看的规则文档是给人看的两者缺一不可。抽检不可省。清洗完别急着送下游随机抽50~100条肉眼过一遍。我每次抽检都能发现至少一个漏网之鱼已经养成习惯了。脚本可重复执行。数据每天都会更新清洗脚本应该设计成可重复跑的。有人用一次性硬编码临时表第二天新数据来了又要重构一次——这是最大的时间浪费。数据清洗不像建模、可视化那么有成就感它枯燥、琐碎、看不见尽头但恰恰是这些“看不见”的活儿决定了整个数据链路的地基稳不稳。我个人的体会是做数据清洗最需要的能力不是代码写得花哨而是对业务的理解、对异常数据的嗅觉以及那点不达目的不罢休的轴劲。一个数据团队前期的清洗功夫下得越深后期的分析、建模就越省心这个道理做久了自然就懂。
返回列表