尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Python Pandas数据清洗实战:构建健壮可复用的预处理流程

Python Pandas数据清洗实战:构建健壮可复用的预处理流程 在实际工作中我们经常需要处理来自不同渠道、格式各异的数据例如用户行为日志、业务系统订单、第三方接口数据等。这些数据在进入分析或处理流程前往往存在字段缺失、格式混乱、类型不匹配等问题直接使用会导致下游任务失败或分析结果失真。数据清洗作为数据预处理的核心环节其目标就是将原始数据转化为高质量、可用的数据。本文将以一个典型的“用户交易记录清洗”场景为例从零开始详细讲解如何设计一个健壮、可复用的数据清洗流程。我们将使用 Python 的 Pandas 库作为核心工具涵盖数据读取、异常值处理、缺失值填充、格式标准化、去重等关键步骤并最终输出一份干净的数据集。无论你是数据分析师、数据工程师还是后端开发这套方法都能帮助你构建可靠的数据预处理管道。1. 理解数据清洗的核心目标与常见问题数据清洗不是简单的“删除脏数据”而是一个有明确目标的系统工程。其核心在于提升数据的“可用性”为后续的分析、建模或系统集成打下坚实基础。1.1 数据清洗的五大核心目标完整性确保数据记录和字段没有缺失。例如用户ID、交易时间等关键字段必须存在。一致性确保数据在其定义的域内保持一致。例如“性别”字段的值只能是“男”、“女”或“未知”不能出现“M”、“F”或数字1、2。准确性数据必须准确反映真实世界实体或事件。例如用户的年龄不应为负数交易金额应在合理范围内。唯一性避免数据集中存在不必要的重复记录。重复数据会扭曲统计结果如计算总销售额时会被重复计算。时效性数据应在其有效期内被处理和使用。对于时间敏感的分析过时的数据需要被识别或排除。1.2 典型“脏数据”场景与影响假设我们收到一份名为raw_transactions.csv的交易数据它可能包含以下问题缺失值user_id或amount字段为空NaN。格式混乱transaction_time字段可能是字符串“2023-01-01”也可能是时间戳“1672531200”甚至混有“01/01/2023”这种格式。异常值amount字段出现负数或极大值如999999age字段为200。不一致性product_category字段中“电子产品”、“电子商品”、“3C”代表同一含义。重复记录完全相同的行出现了多次。如果不对这些问题进行处理直接进行月度销售额统计、用户画像分析或机器学习模型训练得到的结果将是不可靠的甚至会导致错误的业务决策。2. 环境准备与工具选择我们将使用 Python 和 Pandas 库来完成本次数据清洗实战。Pandas 提供了强大的数据结构和函数是进行数据清洗和预处理的行业标准工具之一。2.1 环境搭建首先确保你的 Python 环境已安装必要的库。推荐使用 Anaconda 或通过pip安装。# 使用 pip 安装 pandas 和 numpy pip install pandas numpy为了后续可能的数据可视化或更复杂的操作也可以一并安装matplotlib和scikit-learn。pip install matplotlib scikit-learn2.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‘, periodsn_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(frac0.05).index, ‘user_id‘] np.nan df_raw.loc[df_raw.sample(frac0.03).index, ‘amount‘] np.nan df_raw.loc[df_raw.sample(frac0.02).index, ‘city‘] np.nan # 2. 制造格式混乱的 transaction_time df_raw.loc[df_raw.sample(frac0.1).index, ‘transaction_time‘] pd.to_datetime(df_raw[‘transaction_time‘]).sample(frac0.1).astype(int) // 10**9 # 转为时间戳 df_raw.loc[df_raw.sample(frac0.05).index, ‘transaction_time‘] ‘01/15/2023 14:30‘ # 混入不同格式 # 3. 制造异常值 df_raw.loc[df_raw.sample(frac0.01).index, ‘amount‘] -100 df_raw.loc[df_raw.sample(frac0.005).index, ‘amount‘] 999999 df_raw.loc[df_raw.sample(frac0.01).index, ‘user_age‘] 200 # 4. 制造不一致的分类值 df_raw.loc[df_raw[‘product_category‘] ‘电子产品‘].sample(frac0.3).index, ‘product_category‘] ‘电子商品‘ df_raw.loc[df_raw[‘product_category‘] ‘电子产品‘].sample(frac0.2).index, ‘product_category‘] ‘3C‘ # 5. 制造重复行完全重复 dup_indices df_raw.sample(n20).index df_duplicates df_raw.loc[dup_indices].copy() df_raw pd.concat([df_raw, df_duplicates], ignore_indexTrue) # 6. 打乱数据顺序 df_raw df_raw.sample(frac1, random_state42).reset_index(dropTrue) # 保存到 raw 文件夹 df_raw.to_csv(‘../data/raw/raw_transactions.csv‘, indexFalse, 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, inplaceTrue) 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(methodmethod) 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_formatTrue) 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, indexdf_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(lowercol_min, uppercol_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_columnsNone, 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(subsetsubset_columns, keepkeep) 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_columnsdup_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_okTrue) df.to_csv(output_path, indexFalse, 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 清洗后数据量骤减现象清洗后的数据行数比原始数据少了很多。排查检查缺失值处理策略是否对关键列如user_id使用了过于严格的‘drop‘策略查看handle_missing_values函数的日志确认删除了多少行。检查格式转换pd.to_datetime或pd.to_numeric转换失败时如果设置了errors‘coerce‘失败的值会变成NaT或NaN这些可能在后续步骤中被删除。检查转换失败的数量。检查去重逻辑drop_duplicates的subset参数是否正确是否把本应保留的记录误判为重复建议在每一步清洗后都打印出数据形状的变化 (df.shape)并记录日志。对于关键列优先考虑填充fill而非删除drop。5.2 内存占用过高或处理速度慢现象处理大型数据集如数GB时脚本运行缓慢或内存溢出。排查与优化指定数据类型在pd.read_csv时使用dtype参数指定每列的数据类型避免Pandas自动推断消耗内存。例如对于ID类字段即使全是数字也可以指定为‘str‘或‘category‘。dtype_spec {‘user_id‘: ‘str‘, ‘product_category‘: ‘category‘, ‘city‘: ‘category‘} df pd.read_csv(filepath, dtypedtype_spec)分块处理使用chunksize参数分批读取和处理数据。chunk_iter pd.read_csv(filepath, chunksize50000) 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_indexTrue)使用高效操作避免在DataFrame上使用循环 (forloop)尽量使用Pandas的向量化操作如replace,fillna,clip。5.3 清洗逻辑错误或覆盖原始数据现象清洗结果不符合预期或者不小心修改了原始数据文件。排查与预防保留原始数据始终在副本 (df.copy()) 上进行操作如我们每个清洗函数内所做的那样。单元测试为每个清洗函数编写简单的单元测试使用小的、可控的测试数据验证其行为。版本控制对清洗脚本和重要的中间数据输出进行版本控制如使用Git。配置化将清洗规则如缺失值策略、映射字典、异常值边界提取到外部配置文件如JSON或YAML中使流程更透明、更易调整。5.4 生产环境数据清洗清单将清洗流程应用于生产环境时需要考虑更多因素考量维度学习/开发环境做法生产环境建议数据来源单个静态文件从数据库、数据仓库、消息队列或对象存储如S3定时拉取任务调度手动运行脚本使用 Airflow, Dagster, Prefect 等调度工具编排清洗任务错误处理打印日志人工查看完善的日志记录如写入ELK错误告警邮件、钉钉、企业微信失败重试机制数据质量监控最终人工检查定义数据质量规则如非空率、唯一性、值域范围使用 Great Expectations、Deequ 等工具在清洗前后自动校验并生成报告代码与配置管理脚本放在本地代码仓库管理清洗规则配置化支持不同环境dev/test/prod的不同参数性能与资源单机运行对于超大数据集考虑使用 PySpark、Dask 进行分布式清洗或使用云服务的ETL工具如AWS Glue6. 扩展方向与总结完成基础清洗后数据预处理流程还可以向更深处扩展特征工程基于清洗后的干净数据可以衍生新的特征。例如从transaction_time中提取“小时”、“星期几”、“是否周末”计算用户的“累计交易金额”、“最近一次交易距今天数”等。管道化与自动化将上述清洗步骤封装成一个Scikit-learn的Transformer类可以轻松地嵌入到机器学习管道中实现训练与预测时数据预处理的一致性。数据质量报告在清洗流程的最后自动生成一份数据质量报告包含处理前后的行数对比、各列缺失率变化、异常值处理情况、唯一值分布等便于审计和追溯。增量清洗对于流式或每日增量的数据设计增量清洗策略只处理新增或变化的数据提升效率。数据清洗是数据价值链的起点其质量直接决定了后续所有环节的上限。一个设计良好的清洗流程应该是可配置、可测试、可监控和可追溯的。本文提供的模块化代码和问题排查思路可以作为一个坚实的起点帮助你根据实际业务数据的复杂程度构建起适合自己的、稳健的数据预处理系统。核心在于理解业务明确每一列数据的含义和约束然后有针对性地应用删除、填充、转换或映射策略并始终对处理结果保持验证的习惯。
返回列表