怎样实现批量清洗采集数据

wen 实用脚本 30

本文目录导读:

怎样实现批量清洗采集数据

  1. 核心流程框架
  2. Python 实战代码示例
  3. 常见清洗场景与代码速查
  4. 进阶优化与扩展
  5. 总结建议

实现批量清洗采集数据,通常需要结合自动化的脚本(如Python)和数据处理库(如Pandas),下面我将从流程框架代码示例常用清洗方法进阶优化四个方面来详细讲解。

核心流程框架

一个标准的批量数据清洗 pipeline 通常包含以下步骤:

  1. 数据加载:读取多个数据源文件(CSV、Excel、JSON、数据库等)。
  2. 统一模式:将不同来源的数据字段名、格式统一。
  3. 清洗规则应用:批量执行去重、缺失值处理、格式修正、异常值过滤。
  4. 验证与输出:检查清洗结果,输出到新文件或数据库。

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'] = None
df['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)

处理大规模数据(超出内存)

使用 daskpolars 库进行分布式或惰性计算。

使用 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-profilingsweetviz 生成清洗前后对比报告。

# 清洗前
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. 先小规模测试:先用1-2个文件测试你的清洗函数,确认效果再批量运行。
  2. 保留原始备份:永远不要直接覆盖原始数据,清洗后的文件另存为新目录。
  3. 记录日志:每一次清洗都记录删了多少行、改了哪些列,方便追溯。
  4. 处理异常编码:国内数据常见 gbkgb2312utf-8-sig,读取时指定 encoding 参数。

如果你能提供你的具体数据类型和想要清洗的难点(比如手机号格式、身份证校验、文本乱码等),我可以帮你写出更针对性的代码。

抱歉,评论功能暂时关闭!