Python Pandas数据清洗实战:构建健壮可复用的预处理流程
2026/8/9 9:46:02 网站建设 项目流程

在实际工作中,我们经常需要处理来自不同渠道、格式各异的数据,例如用户行为日志、业务系统订单、第三方接口数据等。这些数据在进入分析或处理流程前,往往存在字段缺失、格式混乱、类型不匹配等问题,直接使用会导致下游任务失败或分析结果失真。数据清洗,作为数据预处理的核心环节,其目标就是将原始数据转化为高质量、可用的数据。本文将以一个典型的“用户交易记录清洗”场景为例,从零开始,详细讲解如何设计一个健壮、可复用的数据清洗流程。我们将使用 Python 的 Pandas 库作为核心工具,涵盖数据读取、异常值处理、缺失值填充、格式标准化、去重等关键步骤,并最终输出一份干净的数据集。无论你是数据分析师、数据工程师还是后端开发,这套方法都能帮助你构建可靠的数据预处理管道。

1. 理解数据清洗的核心目标与常见问题

数据清洗不是简单的“删除脏数据”,而是一个有明确目标的系统工程。其核心在于提升数据的“可用性”,为后续的分析、建模或系统集成打下坚实基础。

1.1 数据清洗的五大核心目标

  1. 完整性:确保数据记录和字段没有缺失。例如,用户ID、交易时间等关键字段必须存在。
  2. 一致性:确保数据在其定义的域内保持一致。例如,“性别”字段的值只能是“男”、“女”或“未知”,不能出现“M”、“F”或数字1、2。
  3. 准确性:数据必须准确反映真实世界实体或事件。例如,用户的年龄不应为负数,交易金额应在合理范围内。
  4. 唯一性:避免数据集中存在不必要的重复记录。重复数据会扭曲统计结果,如计算总销售额时会被重复计算。
  5. 时效性:数据应在其有效期内被处理和使用。对于时间敏感的分析,过时的数据需要被识别或排除。

1.2 典型“脏数据”场景与影响

假设我们收到一份名为raw_transactions.csv的交易数据,它可能包含以下问题:

  • 缺失值user_idamount字段为空(NaN)。
  • 格式混乱transaction_time字段可能是字符串“2023-01-01”,也可能是时间戳“1672531200”,甚至混有“01/01/2023”这种格式。
  • 异常值amount字段出现负数或极大值(如999999),age字段为200。
  • 不一致性product_category字段中,“电子产品”、“电子商品”、“3C”代表同一含义。
  • 重复记录:完全相同的行出现了多次。

如果不对这些问题进行处理,直接进行月度销售额统计、用户画像分析或机器学习模型训练,得到的结果将是不可靠的,甚至会导致错误的业务决策。

2. 环境准备与工具选择

我们将使用 Python 和 Pandas 库来完成本次数据清洗实战。Pandas 提供了强大的数据结构和函数,是进行数据清洗和预处理的行业标准工具之一。

2.1 环境搭建

首先,确保你的 Python 环境已安装必要的库。推荐使用 Anaconda 或通过pip安装。

# 使用 pip 安装 pandas 和 numpy pip install pandas numpy

为了后续可能的数据可视化或更复杂的操作,也可以一并安装matplotlibscikit-learn

pip install matplotlib scikit-learn

2.2 创建项目结构与模拟数据

创建一个清晰的项目目录,有助于管理脚本、数据和输出结果。

data_cleaning_project/ ├── data/ │ ├── raw/ # 存放原始数据 │ │ └── raw_transactions.csv │ └── cleaned/ # 存放清洗后的数据 ├── scripts/ │ └── data_cleaning.py # 主清洗脚本 ├── notebooks/ # (可选) Jupyter Notebook 用于探索 │ └── exploration.ipynb └── README.md

接下来,我们模拟生成一份包含上述“脏数据”特征的CSV文件,用于演示。将以下 Python 代码保存为scripts/generate_raw_data.py并运行。

import pandas as pd import numpy as np # 设置随机种子保证可复现 np.random.seed(42) # 生成1000条模拟数据 n_samples = 1000 data = { ‘transaction_id‘: range(1000, 1000 + n_samples), ‘user_id‘: [f‘USER_{i:04d}‘ for i in range(n_samples)], ‘transaction_time‘: pd.date_range(‘2023-01-01‘, periods=n_samples, freq=‘h‘).strftime(‘%Y-%m-%d %H:%M:%S‘), ‘amount‘: np.round(np.random.uniform(10, 500, n_samples), 2), ‘product_category‘: np.random.choice([‘电子产品‘, ‘服装‘, ‘食品‘, ‘家居‘, ‘图书‘], n_samples), ‘payment_method‘: np.random.choice([‘支付宝‘, ‘微信支付‘, ‘信用卡‘, ‘银行卡‘], n_samples), ‘user_age‘: np.random.randint(18, 70, n_samples), ‘city‘: np.random.choice([‘北京‘, ‘上海‘, ‘广州‘, ‘深圳‘, ‘杭州‘, ‘成都‘], n_samples) } df_raw = pd.DataFrame(data) # 人为注入“脏数据” # 1. 制造缺失值 df_raw.loc[df_raw.sample(frac=0.05).index, ‘user_id‘] = np.nan df_raw.loc[df_raw.sample(frac=0.03).index, ‘amount‘] = np.nan df_raw.loc[df_raw.sample(frac=0.02).index, ‘city‘] = np.nan # 2. 制造格式混乱的 transaction_time df_raw.loc[df_raw.sample(frac=0.1).index, ‘transaction_time‘] = pd.to_datetime(df_raw[‘transaction_time‘]).sample(frac=0.1).astype(int) // 10**9 # 转为时间戳 df_raw.loc[df_raw.sample(frac=0.05).index, ‘transaction_time‘] = ‘01/15/2023 14:30‘ # 混入不同格式 # 3. 制造异常值 df_raw.loc[df_raw.sample(frac=0.01).index, ‘amount‘] = -100 df_raw.loc[df_raw.sample(frac=0.005).index, ‘amount‘] = 999999 df_raw.loc[df_raw.sample(frac=0.01).index, ‘user_age‘] = 200 # 4. 制造不一致的分类值 df_raw.loc[df_raw[‘product_category‘] == ‘电子产品‘].sample(frac=0.3).index, ‘product_category‘] = ‘电子商品‘ df_raw.loc[df_raw[‘product_category‘] == ‘电子产品‘].sample(frac=0.2).index, ‘product_category‘] = ‘3C‘ # 5. 制造重复行(完全重复) dup_indices = df_raw.sample(n=20).index df_duplicates = df_raw.loc[dup_indices].copy() df_raw = pd.concat([df_raw, df_duplicates], ignore_index=True) # 6. 打乱数据顺序 df_raw = df_raw.sample(frac=1, random_state=42).reset_index(drop=True) # 保存到 raw 文件夹 df_raw.to_csv(‘../data/raw/raw_transactions.csv‘, index=False, encoding=‘utf-8-sig‘) print(“原始数据已生成并保存至 data/raw/raw_transactions.csv“) print(f“数据形状: {df_raw.shape}“) print(df_raw.head()) print(“\n脏数据统计:“) print(f“user_id 缺失数: {df_raw[‘user_id‘].isna().sum()}“) print(f“amount 缺失数: {df_raw[‘amount‘].isna().sum()}“) print(f“amount 负值数: {(df_raw[‘amount‘] < 0).sum()}“) print(f“user_age 异常(>100)数: {(df_raw[‘user_age‘] > 100).sum()}“)

运行此脚本后,你将在data/raw/目录下得到一份用于清洗练习的原始数据文件。

3. 构建模块化的数据清洗流程

一个健壮的清洗流程应该是模块化和可配置的。我们将清洗步骤分解为独立的函数,便于测试、复用和维护。创建主清洗脚本scripts/data_cleaning.py

3.1 数据加载与初步探索

首先,编写一个函数来加载数据并快速了解其概况。

import pandas as pd import numpy as np import os def load_and_explore_data(filepath): """ 加载数据并进行初步探索 """ try: # 尝试自动推断分隔符和编码,常见编码有‘utf-8‘, ‘gbk‘, ‘utf-8-sig‘ df = pd.read_csv(filepath, encoding=‘utf-8-sig‘) print(f“成功加载数据,形状: {df.shape}“) except UnicodeDecodeError: try: df = pd.read_csv(filepath, encoding=‘gbk‘) print(f“使用 gbk 编码成功加载数据,形状: {df.shape}“) except Exception as e: print(f“加载数据失败: {e}“) return None # 查看前几行 print(“\n数据前5行:“) print(df.head()) # 查看列信息 print(“\n数据列信息:“) print(df.info()) # 查看基本统计信息(针对数值列) print(“\n数值列基本统计:“) print(df.describe()) # 查看缺失值情况 print(“\n各列缺失值数量:“) print(df.isnull().sum()) return df if __name__ == ‘__main__‘: raw_data_path = ‘../data/raw/raw_transactions.csv‘ df = load_and_explore_data(raw_data_path)

运行此脚本,你会看到数据的维度、各列数据类型、缺失值数量以及数值列的统计信息(均值、标准差、最小最大值等),这有助于制定具体的清洗策略。

3.2 处理缺失值

缺失值的处理需要根据业务逻辑决定,常见方法有删除、填充和插值。

def handle_missing_values(df, strategy_dict): """ 根据策略字典处理缺失值 strategy_dict 格式: {‘column_name‘: {‘method‘: ‘drop‘|‘fill‘|‘ffill‘|‘bfill‘, ‘value‘: fill_value}} ‘drop‘: 删除该列为空的行(谨慎使用,可能丢失大量数据) ‘fill‘: 用指定值填充 ‘ffill‘/‘bfill‘: 用前向/后向填充(适用于时间序列) """ df_cleaned = df.copy() rows_before = df_cleaned.shape[0] for col, config in strategy_dict.items(): if col not in df_cleaned.columns: print(f“警告: 列 {col} 不存在于数据中“) continue method = config.get(‘method‘) fill_value = config.get(‘value‘) if method == ‘drop‘: # 只删除该特定列为空的行 df_cleaned = df_cleaned.dropna(subset=[col]) print(f“列 [{col}] 采用删除缺失值策略,删除了 {rows_before - df_cleaned.shape[0]} 行。“) rows_before = df_cleaned.shape[0] elif method == ‘fill‘: if fill_value is not None: df_cleaned[col].fillna(fill_value, inplace=True) print(f“列 [{col}] 用值 [{fill_value}] 填充了 {df[col].isna().sum()} 个缺失值。“) else: print(f“警告: 列 [{col}] 指定了 ‘fill‘ 策略但未提供 ‘value‘,已跳过。“) elif method in [‘ffill‘, ‘bfill‘]: df_cleaned[col] = df_cleaned[col].fillna(method=method) print(f“列 [{col}] 采用了 [{method}] 填充。“) else: print(f“警告: 列 [{col}] 的策略 [{method}] 不被支持,已跳过。“) print(f“缺失值处理完成。数据形状从 {df.shape} 变为 {df_cleaned.shape}“) return df_cleaned # 定义缺失值处理策略 missing_strategy = { ‘user_id‘: {‘method‘: ‘drop‘}, # 用户ID是关键标识,缺失则删除该行 ‘amount‘: {‘method‘: ‘fill‘, ‘value‘: df[‘amount‘].median()}, # 金额用中位数填充,避免极端值影响 ‘city‘: {‘method‘: ‘fill‘, ‘value‘: ‘未知‘}, # 城市信息用‘未知‘填充 }

3.3 处理异常值与格式标准化

这一步需要将数据转换为一致的格式,并剔除或修正不合理的值。

def standardize_and_handle_outliers(df, config): """ 标准化格式并处理异常值 config 格式示例: { ‘columns‘: { ‘transaction_time‘: {‘dtype‘: ‘datetime‘, ‘format‘: ‘mixed‘}, # mixed表示尝试自动解析 ‘amount‘: {‘dtype‘: ‘float‘, ‘min‘: 0, ‘max‘: 100000, ‘clip‘: True}, # clip为True则将异常值截断到边界 ‘user_age‘: {‘dtype‘: ‘int‘, ‘min‘: 0, ‘max‘: 120, ‘clip‘: False} # clip为False则将异常值设为NaN } } """ df_processed = df.copy() for col, rules in config.get(‘columns‘, {}).items(): if col not in df_processed.columns: continue target_dtype = rules.get(‘dtype‘) # 1. 格式转换 if target_dtype == ‘datetime‘: # Pandas 的 to_datetime 可以处理多种格式,errors=‘coerce‘将解析失败的设为NaT df_processed[col] = pd.to_datetime(df_processed[col], errors=‘coerce‘, infer_datetime_format=True) print(f“列 [{col}] 已转换为 datetime 格式,转换失败数: {df_processed[col].isna().sum() - df[col].isna().sum()}“) elif target_dtype in [‘int‘, ‘float‘]: df_processed[col] = pd.to_numeric(df_processed[col], errors=‘coerce‘) print(f“列 [{col}] 已转换为 {target_dtype} 格式,转换失败数: {df_processed[col].isna().sum() - df[col].isna().sum()}“) # 2. 处理异常值 (针对数值型) if target_dtype in [‘int‘, ‘float‘]: col_min = rules.get(‘min‘) col_max = rules.get(‘max‘) clip = rules.get(‘clip‘, False) if col_min is not None or col_max is not None: outlier_mask = pd.Series(True, index=df_processed.index) if col_min is not None: outlier_mask &= (df_processed[col] < col_min) if col_max is not None: outlier_mask &= (df_processed[col] > col_max) outlier_count = outlier_mask.sum() if clip: # 截断到边界值 df_processed[col] = df_processed[col].clip(lower=col_min, upper=col_max) print(f“列 [{col}] 截断了 {outlier_count} 个超出范围 [{col_min}, {col_max}] 的值。“) else: # 将异常值设为NaN,后续可由缺失值处理策略处理 df_processed.loc[outlier_mask, col] = np.nan print(f“列 [{col}] 将 {outlier_count} 个超出范围 [{col_min}, {col_max}] 的值标记为缺失。“) return df_processed # 定义格式与异常值处理配置 standardize_config = { ‘columns‘: { ‘transaction_time‘: {‘dtype‘: ‘datetime‘}, ‘amount‘: {‘dtype‘: ‘float‘, ‘min‘: 0.01, ‘max‘: 100000, ‘clip‘: True}, # 交易金额最小0.01元,最大10万元,超出则截断 ‘user_age‘: {‘dtype‘: ‘int‘, ‘min‘: 0, ‘max‘: 120, ‘clip‘: False}, # 年龄异常标记为缺失 } }

3.4 处理数据不一致性与重复值

对于分类数据的不一致和重复记录,需要进行映射和去重。

def standardize_categorical_values(df, mapping_dict): """ 根据映射字典标准化分类变量的值 mapping_dict 格式: {‘column_name‘: {‘old_value1‘: ‘new_value1‘, ‘old_value2‘: ‘new_value2‘}} """ df_mapped = df.copy() for col, value_map in mapping_dict.items(): if col in df_mapped.columns: # 使用 replace 进行映射 df_mapped[col] = df_mapped[col].replace(value_map) print(f“列 [{col}] 已完成值映射标准化。“) return df_mapped def remove_duplicates(df, subset_columns=None, keep=‘first‘): """ 基于指定列子集删除重复行 subset_columns: 判断重复依据的列列表,None则考虑所有列 keep: ‘first‘保留第一条,‘last‘保留最后一条,False删除所有重复项 """ df_deduped = df.copy() rows_before = df_deduped.shape[0] df_deduped = df_deduped.drop_duplicates(subset=subset_columns, keep=keep) rows_after = df_deduped.shape[0] removed = rows_before - rows_after print(f“基于列 {subset_columns} 进行去重,删除了 {removed} 条重复记录。“) return df_deduped # 定义分类值映射 category_mapping = { ‘product_category‘: { ‘电子商品‘: ‘电子产品‘, ‘3C‘: ‘电子产品‘, # 可以继续添加其他映射 } } # 假设我们认为 transaction_id, user_id, transaction_time 共同唯一标识一笔交易 dup_subset = [‘transaction_id‘, ‘user_id‘, ‘transaction_time‘]

4. 组装完整清洗管道并验证结果

现在,我们将所有步骤串联起来,形成一个完整的清洗管道,并验证清洗效果。

def run_cleaning_pipeline(raw_data_path, output_path): """ 执行完整的数据清洗管道 """ print(“=== 开始数据清洗管道 ===“) # 1. 加载与探索 df = load_and_explore_data(raw_data_path) if df is None: return # 2. 处理缺失值 print(“\n--- 步骤1: 处理缺失值 ---“) missing_strategy = { ‘user_id‘: {‘method‘: ‘drop‘}, ‘amount‘: {‘method‘: ‘fill‘, ‘value‘: df[‘amount‘].median()}, ‘city‘: {‘method‘: ‘fill‘, ‘value‘: ‘未知‘}, } df = handle_missing_values(df, missing_strategy) # 3. 标准化格式与处理异常值 print(“\n--- 步骤2: 标准化格式与处理异常值 ---“) standardize_config = { ‘columns‘: { ‘transaction_time‘: {‘dtype‘: ‘datetime‘}, ‘amount‘: {‘dtype‘: ‘float‘, ‘min‘: 0.01, ‘max‘: 100000, ‘clip‘: True}, ‘user_age‘: {‘dtype‘: ‘int‘, ‘min‘: 0, ‘max‘: 120, ‘clip‘: False}, } } df = standardize_and_handle_outliers(df, standardize_config) # 由于异常值可能被设为NaN,再次处理缺失值(例如user_age) print(“\n--- 步骤3: 二次处理缺失值(由异常值转换而来)---“) secondary_missing_strategy = { ‘user_age‘: {‘method‘: ‘fill‘, ‘value‘: int(df[‘user_age‘].median())}, } df = handle_missing_values(df, secondary_missing_strategy) # 4. 标准化分类值 print(“\n--- 步骤4: 标准化分类值 ---“) category_mapping = { ‘product_category‘: {‘电子商品‘: ‘电子产品‘, ‘3C‘: ‘电子产品‘}, } df = standardize_categorical_values(df, category_mapping) # 5. 去重 print(“\n--- 步骤5: 删除重复记录 ---“) # 根据业务逻辑选择去重列,这里假设三列共同唯一 dup_subset = [‘transaction_id‘, ‘user_id‘, ‘transaction_time‘] df = remove_duplicates(df, subset_columns=dup_subset, keep=‘first‘) # 6. 最终数据探索与保存 print(“\n=== 清洗完成,最终数据概览 ===") print(f“最终数据形状: {df.shape}“) print(“\n前5行数据:“) print(df.head()) print(“\n各列数据类型:“) print(df.dtypes) print(“\n缺失值检查:“) print(df.isnull().sum().sum()) # 总和应为0 print(“\n‘product_category‘ 唯一值:“) print(df[‘product_category‘].unique()) # 保存清洗后的数据 os.makedirs(os.path.dirname(output_path), exist_ok=True) df.to_csv(output_path, index=False, encoding=‘utf-8-sig‘) print(f“\n清洗后的数据已保存至: {output_path}“) return df if __name__ == ‘__main__‘: raw_path = ‘../data/raw/raw_transactions.csv‘ cleaned_path = ‘../data/cleaned/cleaned_transactions.csv‘ cleaned_df = run_cleaning_pipeline(raw_path, cleaned_path)

运行这个主脚本,控制台会输出每一步的处理日志。清洗完成后,打开data/cleaned/cleaned_transactions.csv,你将得到一份格式统一、无缺失、无异常、无重复的干净数据集。

5. 常见问题排查与最佳实践

在实际项目中,数据清洗过程不会总是一帆风顺。以下是几个常见问题及其排查思路。

5.1 清洗后数据量骤减

  • 现象:清洗后的数据行数比原始数据少了很多。
  • 排查
    1. 检查缺失值处理策略:是否对关键列(如user_id)使用了过于严格的‘drop‘策略?查看handle_missing_values函数的日志,确认删除了多少行。
    2. 检查格式转换pd.to_datetimepd.to_numeric转换失败时,如果设置了errors=‘coerce‘,失败的值会变成NaTNaN,这些可能在后续步骤中被删除。检查转换失败的数量。
    3. 检查去重逻辑drop_duplicatessubset参数是否正确?是否把本应保留的记录误判为重复?
  • 建议:在每一步清洗后,都打印出数据形状的变化 (df.shape),并记录日志。对于关键列,优先考虑填充(fill)而非删除(drop)。

5.2 内存占用过高或处理速度慢

  • 现象:处理大型数据集(如数GB)时,脚本运行缓慢或内存溢出。
  • 排查与优化
    1. 指定数据类型:在pd.read_csv时使用dtype参数指定每列的数据类型,避免Pandas自动推断消耗内存。例如,对于ID类字段,即使全是数字,也可以指定为‘str‘‘category‘
      dtype_spec = {‘user_id‘: ‘str‘, ‘product_category‘: ‘category‘, ‘city‘: ‘category‘} df = pd.read_csv(filepath, dtype=dtype_spec)
    2. 分块处理:使用chunksize参数分批读取和处理数据。
      chunk_iter = pd.read_csv(filepath, chunksize=50000) cleaned_chunks = [] for chunk in chunk_iter: # 对每个chunk应用清洗函数 cleaned_chunk = some_cleaning_function(chunk) cleaned_chunks.append(cleaned_chunk) df_cleaned = pd.concat(cleaned_chunks, ignore_index=True)
    3. 使用高效操作:避免在DataFrame上使用循环 (forloop),尽量使用Pandas的向量化操作(如replace,fillna,clip)。

5.3 清洗逻辑错误或覆盖原始数据

  • 现象:清洗结果不符合预期,或者不小心修改了原始数据文件。
  • 排查与预防
    1. 保留原始数据:始终在副本 (df.copy()) 上进行操作,如我们每个清洗函数内所做的那样。
    2. 单元测试:为每个清洗函数编写简单的单元测试,使用小的、可控的测试数据验证其行为。
    3. 版本控制:对清洗脚本和重要的中间数据输出进行版本控制(如使用Git)。
    4. 配置化:将清洗规则(如缺失值策略、映射字典、异常值边界)提取到外部配置文件(如JSON或YAML)中,使流程更透明、更易调整。

5.4 生产环境数据清洗清单

将清洗流程应用于生产环境时,需要考虑更多因素:

考量维度学习/开发环境做法生产环境建议
数据来源单个静态文件从数据库、数据仓库、消息队列或对象存储(如S3)定时拉取
任务调度手动运行脚本使用 Airflow, Dagster, Prefect 等调度工具编排清洗任务
错误处理打印日志,人工查看完善的日志记录(如写入ELK),错误告警(邮件、钉钉、企业微信),失败重试机制
数据质量监控最终人工检查定义数据质量规则(如非空率、唯一性、值域范围),使用 Great Expectations、Deequ 等工具在清洗前后自动校验并生成报告
代码与配置管理脚本放在本地代码仓库管理,清洗规则配置化,支持不同环境(dev/test/prod)的不同参数
性能与资源单机运行对于超大数据集,考虑使用 PySpark、Dask 进行分布式清洗,或使用云服务的ETL工具(如AWS Glue)

6. 扩展方向与总结

完成基础清洗后,数据预处理流程还可以向更深处扩展:

  1. 特征工程:基于清洗后的干净数据,可以衍生新的特征。例如,从transaction_time中提取“小时”、“星期几”、“是否周末”;计算用户的“累计交易金额”、“最近一次交易距今天数”等。
  2. 管道化与自动化:将上述清洗步骤封装成一个Scikit-learnTransformer类,可以轻松地嵌入到机器学习管道中,实现训练与预测时数据预处理的一致性。
  3. 数据质量报告:在清洗流程的最后,自动生成一份数据质量报告,包含处理前后的行数对比、各列缺失率变化、异常值处理情况、唯一值分布等,便于审计和追溯。
  4. 增量清洗:对于流式或每日增量的数据,设计增量清洗策略,只处理新增或变化的数据,提升效率。

数据清洗是数据价值链的起点,其质量直接决定了后续所有环节的上限。一个设计良好的清洗流程应该是可配置、可测试、可监控和可追溯的。本文提供的模块化代码和问题排查思路,可以作为一个坚实的起点,帮助你根据实际业务数据的复杂程度,构建起适合自己的、稳健的数据预处理系统。核心在于理解业务,明确每一列数据的含义和约束,然后有针对性地应用删除、填充、转换或映射策略,并始终对处理结果保持验证的习惯。

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

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

立即咨询