本文目录导读:

实现批量清洗采集数据,通常需要结合自动化的脚本(如Python)和数据处理库(如Pandas),下面我将从流程框架、代码示例、常用清洗方法和进阶优化四个方面来详细讲解。
核心流程框架
一个标准的批量数据清洗 pipeline 通常包含以下步骤:
- 数据加载:读取多个数据源文件(CSV、Excel、JSON、数据库等)。
- 统一模式:将不同来源的数据字段名、格式统一。
- 清洗规则应用:批量执行去重、缺失值处理、格式修正、异常值过滤。
- 验证与输出:检查清洗结果,输出到新文件或数据库。
Python 实战代码示例
假设你有多个 CSV 文件存储在 ./raw_data/ 文件夹下,你需要批量清洗它们。
准备工作
import pandas as pd import os from datetime import datetime import re # 配置路径 INPUT_DIR = "./raw_data/" OUTPUT_DIR = "./cleaned_data/" LOG_FILE = "./cleaning_log.txt" # 确保输出目录存在 os.makedirs(OUTPUT_DIR, exist_ok=True)
定义清洗函数(核心)
def clean_dataframe(df, filename):
"""
对单个DataFrame执行清洗,返回清洗后的数据和日志。
"""
original_count = len(df)
log_messages = []
# --- 清洗规则 1:去除完全重复的行 ---
df.drop_duplicates(inplace=True)
dup_count = original_count - len(df)
if dup_count > 0:
log_messages.append(f"去除了 {dup_count} 行重复数据。")
# --- 清洗规则 2:统一列名(去掉空格、特殊字符,转小写) ---
df.columns = [re.sub(r'[^a-zA-Z0-9_]', '_', col).strip().lower() for col in df.columns]
# --- 清洗规则 3:处理缺失值 ---
# 缺失率大于50%的列直接删除
threshold = 0.5
df.dropna(thresh=len(df) * (1 - threshold), axis=1, inplace=True)
# 对于数值列,用中位数填充
numeric_cols = df.select_dtypes(include=['number']).columns
df[numeric_cols] = df[numeric_cols].fillna(df[numeric_cols].median())
# 对于文本列,用 '未知' 填充
text_cols = df.select_dtypes(include=['object']).columns
df[text_cols] = df[text_cols].fillna('未知')
# --- 清洗规则 4:数据类型转换 ---
# 尝试将 'date' 列转为日期类型
if 'date' in df.columns:
df['date'] = pd.to_datetime(df['date'], errors='coerce')
# 转换失败的行保留为 NaT,后续处理
# --- 清洗规则 5:文本标准化(去除前后空格,统一大小写) ---
for col in text_cols:
if df[col].dtype == 'object':
df[col] = df[col].str.strip().str.lower()
# --- 清洗规则 6:异常值过滤(例如年龄在0-120之外) ---
if 'age' in df.columns:
df = df[(df['age'] >= 0) & (df['age'] <= 120)]
# 日志记录
final_count = len(df)
log_messages.append(f"原始行数: {original_count}, 清洗后行数: {final_count}, 删除: {original_count - final_count} 行。")
return df, log_messages
批量处理所有文件
def batch_clean(input_dir, output_dir, log_file):
all_logs = []
# 遍历所有 CSV 文件
for filename in os.listdir(input_dir):
if filename.endswith('.csv'):
filepath = os.path.join(input_dir, filename)
print(f"正在清洗: {filename}")
try:
# 读取原始数据(你可以指定编码解决中文乱码)
df = pd.read_csv(filepath, encoding='utf-8-sig')
# 执行清洗
cleaned_df, logs = clean_dataframe(df, filename)
# 保存清洗后的文件
output_path = os.path.join(output_dir, f"cleaned_{filename}")
cleaned_df.to_csv(output_path, index=False, encoding='utf-8-sig')
# 记录日志
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
for log in logs:
all_logs.append(f"[{timestamp}] {filename}: {log}")
except Exception as e:
error_msg = f"[ERROR] {filename}: 清洗失败 - {str(e)}"
print(error_msg)
all_logs.append(error_msg)
# 写入日志文件
with open(log_file, 'w', encoding='utf-8') as f:
f.write("===== 数据清洗日志 =====\n")
f.write(f"总文件数: {len([f for f in os.listdir(input_dir) if f.endswith('.csv')])}\n\n")
for line in all_logs:
f.write(line + "\n")
if __name__ == "__main__":
batch_clean(INPUT_DIR, OUTPUT_DIR, LOG_FILE)
print("批量清洗完成!")
常见清洗场景与代码速查
| 需要处理的问题 | 代码实现 (Pandas) |
|---|---|
| 去除空格/特殊字符 | df['col'] = df['col'].str.strip()df['col'] = df['col'].str.replace(r'[^\w]', '', regex=True) |
| 统一日期格式 | df['date'] = pd.to_datetime(df['date'], format='%Y-%m-%d') |
| 处理错误数据 | df.loc[df['price'] < 0, 'price'] = Nonedf['price'].fillna(df['price'].median()) |
| 英文大小写统一 | df['name'] = df['name'].str.lower() |
| 按条件删除行 | df = df[df['status'].isin(['active', 'pending'])] |
| 保留特定列 | df = df[['col1', 'col2', 'col3']] |
| 合并多个来源 | combined_df = pd.concat([df1, df2], ignore_index=True) |
进阶优化与扩展
使用配置文件管理规则
将清洗规则写在 YAML 或 JSON 文件中,方便非技术人员修改。
# rules.yaml
cleaning:
deduplicate: true
fill_na:
numeric: median
categorical: "未知"
remove_outliers:
columns: ["age", "price"]
method: iqr # 或 z-score
在 Python 中读取:
import yaml
with open('rules.yaml') as f:
rules = yaml.safe_load(f)
处理大规模数据(超出内存)
使用 dask 或 polars 库进行分布式或惰性计算。
使用 Dask(类似 Pandas 但支持并行):
import dask.dataframe as dd
# 读取多个文件(支持通配符)
ddf = dd.read_csv("./raw_data/*.csv", dtype={'col': 'str'})
# 应用清洗函数
ddf_cleaned = ddf.map_partitions(clean_dataframe)
# 输出到单个文件或分区
ddf_cleaned.to_csv("./cleaned_data/*.csv", index=False)
构建可视化清洗报告
使用 pandas-profiling 或 sweetviz 生成清洗前后对比报告。
# 清洗前
import sweetviz as sv
report = sv.analyze(df)
report.show_html("report_before.html")
# 清洗后
report_clean = sv.analyze(cleaned_df)
report_clean.show_html("report_after.html")
总结建议
- 先小规模测试:先用1-2个文件测试你的清洗函数,确认效果再批量运行。
- 保留原始备份:永远不要直接覆盖原始数据,清洗后的文件另存为新目录。
- 记录日志:每一次清洗都记录删了多少行、改了哪些列,方便追溯。
- 处理异常编码:国内数据常见
gbk、gb2312、utf-8-sig,读取时指定encoding参数。
如果你能提供你的具体数据类型和想要清洗的难点(比如手机号格式、身份证校验、文本乱码等),我可以帮你写出更针对性的代码。