本文目录导读:

我来详细介绍如何配置数据清洗规则,涵盖不同场景和实现方式:
配置化清洗规则设计
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)
最佳实践建议
- 规则验证:添加规则验证逻辑
- 日志记录:记录每次清洗的变更
- 可追溯性:保留原始数据快照
- 渐进式应用:允许分批应用规则
- 性能优化:支持并行处理和批量操作
- 版本控制:管理规则版本历史
这样的配置化设计使得数据清洗规则可以灵活管理,易于维护和扩展。