脚本如何配置数据清洗规则

wen 实用脚本 29

本文目录导读:

脚本如何配置数据清洗规则

  1. 配置化清洗规则设计
  2. Python 实现示例
  3. 数据库配置方式
  4. Web 界面配置示例
  5. 最佳实践建议

我来详细介绍如何配置数据清洗规则,涵盖不同场景和实现方式:

配置化清洗规则设计

YAML 配置文件格式

# data_cleaning_rules.yaml
rules:
  # 缺失值处理
  - field: "age"
    type: "missing_value"
    action: "fill"
    fill_value: 0
    condition: "isnull"
  - field: "email"
    type: "missing_value"
    action: "drop_row"
    condition: "isnull"
  - field: "salary"
    type: "missing_value"
    action: "fill_with_mean"
    group_by: "department"
  # 异常值处理
  - field: "age"
    type: "outlier"
    method: "z_score"
    threshold: 3
    action: "cap"  # 或 "remove"
  - field: "price"
    type: "outlier"
    method: "iqr"
    action: "replace_with_median"
  # 数据类型转换
  - field: "date"
    type: "type_conversion"
    target_type: "datetime"
    format: "%Y-%m-%d"
  - field: "phone"
    type: "type_conversion"
    target_type: "string"
  # 格式标准化
  - field: "name"
    type: "normalization"
    rules:
      - trim: true
      - lowercase: false
      - remove_special_chars: true
  - field: "phone"
    type: "formatting"
    pattern: "^(\\d{3})-\\d{4}-\\d{4}$"
    replacement: "$1-****-$2"
  # 重复数据
  - type: "deduplication"
    subset: ["email", "phone"]
    keep: "first"
  # 自定义规则
  - field: "score"
    type: "custom"
    validation: "value >= 0 and value <= 100"

JSON 配置格式

{
  "cleaning_rules": [
    {
      "field": "age",
      "rules": [
        {"type": "missing", "action": "fill", "value": 0},
        {"type": "range", "min": 0, "max": 150, "action": "clip"},
        {"type": "type", "target": "int"}
      ]
    },
    {
      "field": "salary",
      "rules": [
        {"type": "missing", "action": "mean"},
        {"type": "outlier", "method": "iqr", "action": "median"}
      ]
    }
  ]
}

Python 实现示例

核心清洗引擎

import pandas as pd
import numpy as np
from scipy import stats
import yaml
import json
import re
class DataCleaner:
    def __init__(self, config_path):
        with open(config_path, 'r') as f:
            if config_path.endswith('.yaml'):
                self.config = yaml.safe_load(f)
            else:
                self.config = json.load(f)
        self.rule_handlers = {
            'missing_value': self.handle_missing,
            'outlier': self.handle_outlier,
            'type_conversion': self.handle_type_conversion,
            'normalization': self.handle_normalization,
            'formatting': self.handle_formatting,
            'deduplication': self.handle_deduplication,
            'custom': self.handle_custom
        }
    def clean(self, df):
        """执行所有清洗规则"""
        cleaned_df = df.copy()
        rules = self.config.get('rules', [])
        for rule in rules:
            rule_type = rule.get('type')
            handler = self.rule_handlers.get(rule_type)
            if handler:
                try:
                    cleaned_df = handler(cleaned_df, rule)
                    print(f"Rule {rule_type} applied successfully")
                except Exception as e:
                    print(f"Error applying rule {rule_type}: {e}")
        return cleaned_df
    def handle_missing(self, df, rule):
        """处理缺失值"""
        field = rule.get('field')
        action = rule.get('action')
        if action == 'drop_row':
            df = df.dropna(subset=[field])
        elif action == 'fill':
            df[field] = df[field].fillna(rule.get('fill_value', 0))
        elif action == 'fill_with_mean':
            group_by = rule.get('group_by')
            if group_by:
                df[field] = df.groupby(group_by)[field].transform(
                    lambda x: x.fillna(x.mean())
                )
            else:
                df[field] = df[field].fillna(df[field].mean())
        elif action == 'fill_with_median':
            df[field] = df[field].fillna(df[field].median())
        elif action == 'fill_with_mode':
            df[field] = df[field].fillna(df[field].mode()[0])
        elif action == 'forward_fill':
            df[field] = df[field].fillna(method='ffill')
        elif action == 'backward_fill':
            df[field] = df[field].fillna(method='bfill')
        elif action == 'interpolate':
            df[field] = df[field].interpolate()
        return df
    def handle_outlier(self, df, rule):
        """处理异常值"""
        field = rule.get('field')
        method = rule.get('method')
        action = rule.get('action')
        threshold = rule.get('threshold', 3)
        if method == 'z_score':
            z_scores = np.abs(stats.zscore(df[field].dropna()))
            outliers = z_scores > threshold
        elif method == 'iqr':
            Q1 = df[field].quantile(0.25)
            Q3 = df[field].quantile(0.75)
            IQR = Q3 - Q1
            lower_bound = Q1 - 1.5 * IQR
            upper_bound = Q3 + 1.5 * IQR
            outliers = (df[field] < lower_bound) | (df[field] > upper_bound)
        elif method == 'mad':
            median = df[field].median()
            mad = np.median(np.abs(df[field] - median))
            modified_z_score = 0.6745 * (df[field] - median) / mad
            outliers = np.abs(modified_z_score) > threshold
        if action == 'remove':
            df = df[~outliers]
        elif action == 'cap':
            if method == 'z_score':
                df.loc[outliers, field] = df[field].median()
            elif method == 'iqr':
                df.loc[df[field] < lower_bound, field] = lower_bound
                df.loc[df[field] > upper_bound, field] = upper_bound
        elif action == 'replace_with_mean':
            df.loc[outliers, field] = df[field].mean()
        elif action == 'replace_with_median':
            df.loc[outliers, field] = df[field].median()
        return df
    def handle_type_conversion(self, df, rule):
        """数据类型转换"""
        field = rule.get('field')
        target_type = rule.get('target_type')
        format_str = rule.get('format')
        if target_type == 'datetime' and format_str:
            df[field] = pd.to_datetime(df[field], format=format_str, errors='coerce')
        elif target_type == 'int':
            df[field] = pd.to_numeric(df[field], errors='coerce').astype('Int64')
        elif target_type == 'float':
            df[field] = pd.to_numeric(df[field], errors='coerce')
        elif target_type == 'string':
            df[field] = df[field].astype(str)
        elif target_type == 'bool':
            df[field] = df[field].astype(bool)
        elif target_type == 'category':
            df[field] = df[field].astype('category')
        return df
    def handle_normalization(self, df, rule):
        """数据标准化"""
        field = rule.get('field')
        rules = rule.get('rules', {})
        config = {k: v for rule_list in rules 
                  for k, v in rule_list.items()}
        if config.get('trim', False):
            df[field] = df[field].str.strip()
        if config.get('lowercase', False):
            df[field] = df[field].str.lower()
        if config.get('uppercase', False):
            df[field] = df[field].str.upper()
        if config.get('remove_special_chars', False):
            df[field] = df[field].str.replace(r'[^\w\s]', '', regex=True)
        if config.get('remove_whitespace', False):
            df[field] = df[field].str.replace(r'\s+', ' ', regex=True)
        return df
    def handle_formatting(self, df, rule):
        """格式规范化"""
        field = rule.get('field')
        pattern = rule.get('pattern')
        replacement = rule.get('replacement', '')
        df[field] = df[field].str.replace(pattern, replacement, regex=True)
        return df
    def handle_deduplication(self, df, rule):
        """去重处理"""
        subset = rule.get('subset')
        keep = rule.get('keep', 'first')
        if subset:
            df = df.drop_duplicates(subset=subset, keep=keep)
        else:
            df = df.drop_duplicates(keep=keep)
        return df
    def handle_custom(self, df, rule):
        """自定义规则"""
        field = rule.get('field')
        validation = rule.get('validation')
        if validation:
            mask = df[field].apply(lambda x: eval(validation))
            if rule.get('action') == 'remove':
                df = df[mask]
            else:
                df.loc[~mask, field] = rule.get('replacement', np.nan)
        return df
# 使用示例
cleaner = DataCleaner('data_cleaning_rules.yaml')
cleaned_data = cleaner.clean(df)

更灵活的规则配置系统

from typing import List, Dict, Any, Callable
from dataclasses import dataclass
from enum import Enum
class RuleType(Enum):
    MISSING_VALUE = "missing_value"
    OUTLIER = "outlier"
    TYPE_CONVERSION = "type_conversion"
    NORMALIZATION = "normalization"
    FORMATTING = "formatting"
    DEDUPLICATION = "deduplication"
    VALIDATION = "validation"
    CUSTOM = "custom"
@dataclass
class Rule:
    type: RuleType
    field: str = None
    action: str = None
    condition: str = None
    params: Dict[str, Any] = None
    custom_func: Callable = None
class AdvancedDataCleaner:
    def __init__(self, rules: List[Rule] = None):
        self.rules = rules or []
        self.rule_executors = {
            RuleType.MISSING_VALUE: self._handle_missing_value,
            RuleType.OUTLIER: self._handle_outlier,
            RuleType.TYPE_CONVERSION: self._handle_type_conversion,
            RuleType.NORMALIZATION: self._handle_normalization,
            RuleType.FORMATTING: self._handle_formatting,
            RuleType.DEDUPLICATION: self._handle_deduplication,
            RuleType.VALIDATION: self._handle_validation,
            RuleType.CUSTOM: self._handle_custom
        }
    def add_rule(self, rule: Rule):
        self.rules.append(rule)
    def add_rules_from_dict(self, rules_dict: List[Dict]):
        for rule_dict in rules_dict:
            rule = Rule(
                type=RuleType(rule_dict['type']),
                field=rule_dict.get('field'),
                action=rule_dict.get('action'),
                condition=rule_dict.get('condition'),
                params=rule_dict.get('params', {})
            )
            self.add_rule(rule)
    def clean(self, df):
        """执行所有清洗规则"""
        result_df = df.copy()
        for rule in self.rules:
            executor = self.rule_executors.get(rule.type)
            if executor:
                try:
                    result_df = executor(result_df, rule)
                except Exception as e:
                    print(f"Error executing rule {rule.type}: {e}")
        return result_df
    def _handle_missing_value(self, df, rule):
        """处理缺失值"""
        if rule.condition:
            mask = df[rule.field].apply(
                lambda x: eval(rule.condition)
            )
        else:
            mask = df[rule.field].isna()
        if rule.action == 'drop':
            df = df[~mask]
        elif rule.action == 'fill':
            fill_value = rule.params.get('fill_value', 0)
            df.loc[mask, rule.field] = fill_value
        elif rule.action == 'fill_with':
            fill_method = rule.params.get('method', 'mean')
            if fill_method == 'mean':
                df[rule.field] = df[rule.field].fillna(
                    df[rule.field].mean()
                )
            elif fill_method == 'median':
                df[rule.field] = df[rule.field].fillna(
                    df[rule.field].median()
                )
            elif fill_method == 'mode':
                df[rule.field] = df[rule.field].fillna(
                    df[rule.field].mode().iloc[0]
                )
        return df
    def _handle_custom(self, df, rule):
        """处理自定义规则"""
        if rule.custom_func:
            return rule.custom_func(df, rule)
        return df

数据库配置方式

SQL 配置表

-- 创建清洗规则表
CREATE TABLE data_cleaning_rules (
    id INT PRIMARY KEY AUTO_INCREMENT,
    rule_name VARCHAR(100),
    table_name VARCHAR(100),
    field_name VARCHAR(100),
    rule_type VARCHAR(50),
    action VARCHAR(50),
    parameters JSON,
    enabled BOOLEAN DEFAULT TRUE,
    priority INT DEFAULT 0,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
-- 插入规则
INSERT INTO data_cleaning_rules 
    (rule_name, table_name, field_name, rule_type, action, parameters)
VALUES
    ('handle_missing_age', 'users', 'age', 'missing_value', 'fill', 
     JSON_OBJECT('fill_value', 0, 'condition', 'isnull')),
    ('remove_outlier_salary', 'employees', 'salary', 'outlier', 'cap',
     JSON_OBJECT('method', 'iqr', 'threshold', 1.5));

Python SQLite 配置示例

import sqlite3
import json
class DatabaseConfigCleaner:
    def __init__(self, db_path):
        self.conn = sqlite3.connect(db_path)
        self.create_tables()
    def create_tables(self):
        cursor = self.conn.cursor()
        cursor.execute('''
            CREATE TABLE IF NOT EXISTS cleaning_rules (
                id INTEGER PRIMARY KEY AUTOINCREMENT,
                rule_name TEXT,
                table_name TEXT,
                field_name TEXT,
                rule_type TEXT,
                action TEXT,
                parameters TEXT,
                enabled BOOLEAN DEFAULT 1
            )
        ''')
        self.conn.commit()
    def load_rules(self, table_name=None):
        cursor = self.conn.cursor()
        if table_name:
            cursor.execute(
                "SELECT * FROM cleaning_rules WHERE table_name=? AND enabled=1",
                (table_name,)
            )
        else:
            cursor.execute(
                "SELECT * FROM cleaning_rules WHERE enabled=1"
            )
        rules = []
        for row in cursor.fetchall():
            rule = {
                'id': row[0],
                'rule_name': row[1],
                'table_name': row[2],
                'field_name': row[3],
                'rule_type': row[4],
                'action': row[5],
                'parameters': json.loads(row[6])
            }
            rules.append(rule)
        return rules

Web 界面配置示例

# Flask 示例
from flask import Flask, request, jsonify
import pandas as pd
class CleaningRuleAPI:
    def __init__(self):
        self.app = Flask(__name__)
        self.setup_routes()
    def setup_routes(self):
        @self.app.route('/api/rules', methods=['POST'])
        def create_rule():
            data = request.json
            # 验证规则格式
            if self.validate_rule(data):
                return jsonify({'status': 'success', 'rule': data}), 201
            return jsonify({'status': 'error', 'message': 'Invalid rule format'}), 400
        @self.app.route('/api/clean', methods=['POST'])
        def apply_cleaning():
            data = request.json
            rules = data.get('rules', [])
            # 应用清洗规则
            result = self.execute_cleaning(rules)
            return jsonify(result)

最佳实践建议

  1. 规则验证:添加规则验证逻辑
  2. 日志记录:记录每次清洗的变更
  3. 可追溯性:保留原始数据快照
  4. 渐进式应用:允许分批应用规则
  5. 性能优化:支持并行处理和批量操作
  6. 版本控制:管理规则版本历史

这样的配置化设计使得数据清洗规则可以灵活管理,易于维护和扩展。

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