数据清洗是数据分析和脚本开发中绕不开的环节。原始表里常见重复记录、缺失值、格式错误,甚至明显违背业务常识的条目。临时写 drop_duplicates() 或 fillna() 能救急,但换一份数据又要重写。更稳妥的方式是用 Pandas 处理表格,用 Pydantic 定义有效规则,再把清洗、验证、报告串成一条可复用的数据管道。
依赖安装很简单:- pip install pandas pydantic
复制代码
一、先定义什么算有效数据
Pydantic 的 BaseModel 适合描述业务规则。以社区团购订单为例,字段包括姓名、年龄、手机号和订单金额。年龄限制在 0 到 100 之间,手机号必须是 11 位数字且以 1 开头,金额不能为负数。注意 @field_validator 需要和 @classmethod 搭配使用。
- from pydantic import BaseModel, field_validator, ValidationError
- from typing import Optional
- import pandas as pd
- import numpy as np
- class OrderValidator(BaseModel):
- name: str
- age: Optional[int] = None
- phone: Optional[str] = None
- amount: Optional[float] = None
- @field_validator('age')
- @classmethod
- def validate_age(cls, v):
- if v is not None and (v < 0 or v > 100):
- raise ValueError('年龄必须在0到100之间')
- return v
- @field_validator('phone')
- @classmethod
- def validate_phone(cls, v):
- if v and (not v.isdigit() or len(v) != 11 or not v.startswith('1')):
- raise ValueError('手机号必须是11位数字且以1开头')
- return v
- @field_validator('amount')
- @classmethod
- def validate_amount(cls, v):
- if v is not None and v < 0:
- raise ValueError('金额不能为负数')
- return v
复制代码
二、管道骨架与统计记账
定义好验证模式后,创建 DataPipeline。构造函数中初始化 cleaning_stats,记录 duplicates_removed、nulls_handled 和 validation_errors。这样每次运行后都能知道去掉了多少重复、补了多少空、拦下了多少错误。
- class DataPipeline:
- def __init__(self):
- self.cleaning_stats = {
- 'duplicates_removed': 0,
- 'nulls_handled': 0,
- 'validation_errors': 0
- }
复制代码
三、清洗:先去重,再补缺
clean_data 接收 DataFrame,返回处理后的版本。第一步是 drop_duplicates(),而且必须在填充缺失值之前执行。原因是重复记录会影响中位数计算,如果先填充再取中位数,统计基准会被重复行拉偏。去重之后,数值列用中位数填充,因为中位数比均值更抗离群值;文本列统一填成 Unknown,方便后续人工排查。每处理一处空值,nulls_handled 就累加对应数量。
- def clean_data(self, df: pd.DataFrame) -> pd.DataFrame:
- initial_rows = len(df)
- df = df.drop_duplicates()
- self.cleaning_stats['duplicates_removed'] = initial_rows - len(df)
- numeric_columns = df.select_dtypes(include=[np.number]).columns
- for col in numeric_columns:
- null_count = df[col].isnull().sum()
- if null_count > 0:
- df[col] = df[col].fillna(df[col].median())
- self.cleaning_stats['nulls_handled'] += null_count
- string_columns = df.select_dtypes(include=['object']).columns
- for col in string_columns:
- null_count = df[col].isnull().sum()
- if null_count > 0:
- df[col] = df[col].fillna('Unknown')
- self.cleaning_stats['nulls_handled'] += null_count
- return df
复制代码
四、验证:逐行检查,坏数据不中断管道
清洗后的数据看起来完整,但不一定正确。validate_data 会逐行迭代,把每一行转成字典后传给 OrderValidator。有效行收集到 valid_rows,无效行捕获 ValidationError,并记录行号和错误详情。这样单条坏数据不会让整个管道崩溃,能处理的尽量处理,问题则记录在案供人工审查。
- def validate_data(self, df: pd.DataFrame):
- valid_rows = []
- errors = []
- for idx, row in df.iterrows():
- try:
- validated = OrderValidator(**row.to_dict())
- valid_rows.append(validated.model_dump())
- except ValidationError as e:
- errors.append({'row': idx, 'errors': str(e)})
- self.cleaning_stats['validation_errors'] = len(errors)
- return pd.DataFrame(valid_rows), errors
复制代码
五、编排:把清洗和验证串起来
process 方法作为统一入口,先用 df.copy() 调用 clean_data,再把清洗结果交给 validate_data,最后返回 cleaned_data、validation_errors 和 stats 的综合报告。调用方只需要一个入口就能拿到全部结果。
- def process(self, df: pd.DataFrame):
- cleaned_df = self.clean_data(df.copy())
- validated_df, validation_errors = self.validate_data(cleaned_df)
- return {
- 'cleaned_data': validated_df,
- 'validation_errors': validation_errors,
- 'stats': self.cleaning_stats
- }
复制代码
六、社区团购订单场景实测
社区团购订单常来自微信群接龙、小程序下单或团长手动录入,容易出现同一用户重复下单、手机号少填、年龄乱填、金额为负或空白等情况。下面构造一份示例数据:
- if __name__ == '__main__':
- sample_data = pd.DataFrame(
- {
- 'name': ['张伟', '李娜', '王芳', None, '刘洋', '李娜'],
- 'age': [28, -5, 35, 42, 150, -5],
- 'phone': [
- '13812345678',
- '12345',
- '13987654321',
- '13712345678',
- '13612345678',
- '12345',
- ],
- 'amount': [45.5, 60.0, 45.5, None, 88.8, 60.0],
- }
- )
- print('=' * 60)
- print('原始数据:')
- print(sample_data)
- print('=' * 60)
- pipeline = DataPipeline()
- result = pipeline.process(sample_data)
- print()
- print('清洗后的有效数据:')
- print(result['cleaned_data'])
- print()
- print('验证错误清单:')
- if result['validation_errors']:
- for err in result['validation_errors']:
- print(' 行号 {}:{}'.format(err['row'], err['errors']))
- else:
- print(' 无')
- print()
- print('统计信息:')
- for key, value in result['stats'].items():
- print(' {}: {}'.format(key, value))
- print()
- print('=' * 60)
- print(
- '处理完成。有效数据 {} 行,错误 {} 行。'.format(
- len(result['cleaned_data']), len(result['validation_errors'])
- )
- )
复制代码
运行后,管道会先去掉一条完全重复的“李娜”记录,把缺失的姓名填成 Unknown,用中位数填补缺失的金额。接着逐行验证:李娜的年龄 -5 和手机号 12345 会被拦下,刘洋的年龄 150 也会被标记为错误。最终有效数据中,张伟、王芳以及姓名被补为 Unknown 的用户通过;错误清单列出问题行号和原因。团长拿到报告后,可以直接定位哪些订单需要核对手机号,哪些年龄明显是误填。
七、可扩展方向
这条管道核心代码不到 100 行,但已经具备可维护系统的雏形。可以增加自定义清洗规则,例如标准化地址、脱敏手机号;也可以让 Pydantic 模式可配置,使同一管道适配不同来源的数据;数据量很大时,还能把逐行验证换成向量化操作或并行处理。重点不是一次写得多完美,而是拥有可复用骨架,把重复清洗和验证交给管道,把精力放在从干净数据中提取洞察上。
通过 Pandas 负责表格操作、Pydantic 负责规则校验、DataPipeline 负责编排和统计,一条可复用的数据清洗与验证管道就完成了。 |