一份可直接使用的工业级数据清洗模板,覆盖了金融风控场景需求
一、完整的数据清洗Pipeline
"""金融风控数据清洗完整流程作者:智能风控实战专家适用场景:信贷审批、反欺诈、客户评分等"""import pandas as pdimport numpy as npfrom datetime import datetime, timedeltaimport warningswarnings.filterwarnings('ignore')classRiskDataCleaner:"""金融风控数据清洗器"""def__init__(self, config=None):""" 初始化清洗器 Parameters: ----------- config : dict, 可选 配置参数,包括: - missing_threshold: 缺失率阈值(默认0.5) - outlier_method: 异常值处理方法(默认'IQR') - datetime_format: 日期格式 """self.config = {'missing_threshold': 0.5,'outlier_method': 'IQR', # 'IQR'或'3sigma''datetime_format': '%Y-%m-%d','min_age': 18,'max_age': 70,'min_income': 1000,'max_income': 1000000,'save_log': True,'log_file': 'data_cleaning_log.txt' }if config:self.config.update(config)self.cleaning_log = []self.feature_stats = {}deflog_operation(self, operation, details):"""记录清洗操作""" timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S') log_entry = f"[{timestamp}] {operation}: {details}"self.cleaning_log.append(log_entry)print(log_entry)defsave_log(self):"""保存清洗日志"""ifself.config['save_log']:withopen(self.config['log_file'], 'w', encoding='utf-8') as f:for log inself.cleaning_log: f.write(log + '\n')print(f"清洗日志已保存到: {self.config['log_file']}")defload_data(self, file_path, file_type='csv', **kwargs):""" 加载数据 Parameters: ----------- file_path : str 文件路径 file_type : str 文件类型:'csv', 'excel', 'parquet' **kwargs : dict pandas读取参数 Returns: -------- pd.DataFrame 加载的数据 """self.log_operation("数据加载", f"从 {file_path} 加载{file_type}数据")try:if file_type == 'csv': df = pd.read_csv(file_path, **kwargs)elif file_type == 'excel': df = pd.read_excel(file_path, **kwargs)elif file_type == 'parquet': df = pd.read_parquet(file_path, **kwargs)else:raise ValueError(f"不支持的文件类型: {file_type}")self.log_operation("数据加载", f"成功加载数据,形状: {df.shape}")return dfexcept Exception as e:self.log_operation("数据加载", f"加载失败: {str(e)}")raisedefanalyze_data_quality(self, df):""" 分析数据质量 Parameters: ----------- df : pd.DataFrame 原始数据 Returns: -------- dict 数据质量报告 """self.log_operation("数据质量分析", "开始数据质量分析") total_rows = len(df) total_cols = len(df.columns)# 计算缺失值 missing_stats = df.isnull().sum() missing_pct = (missing_stats / total_rows * 100).round(2)# 数据类型分布 dtype_counts = df.dtypes.value_counts().to_dict()# 数据质量报告 quality_report = {'数据形状': f"{total_rows} 行 × {total_cols} 列",'缺失值统计': {'总缺失值': int(missing_stats.sum()),'缺失率最高的5个特征': missing_pct.sort_values(ascending=False).head(5).to_dict(),'缺失率>50%的特征': list(missing_pct[missing_pct > 50].index) },'数据类型分布': dtype_counts,'内存使用': f"{df.memory_usage(deep=True).sum() / 1024**2:.2f} MB" }# 记录到日志self.log_operation("数据质量分析", f"数据形状: {quality_report['数据形状']}")self.log_operation("数据质量分析", f"总缺失值: {quality_report['缺失值统计']['总缺失值']}")return quality_reportdefhandle_missing_values(self, df, strategy='auto'):""" 处理缺失值 Parameters: ----------- df : pd.DataFrame 原始数据 strategy : str 处理策略:'auto'(自动),'drop'(删除),'fill'(填充) Returns: -------- pd.DataFrame 处理后的数据 """self.log_operation("缺失值处理", f"使用策略: {strategy}") df_clean = df.copy() initial_shape = df_clean.shapeif strategy == 'auto':# 自动策略:高缺失率字段删除,其他填充 missing_pct = df_clean.isnull().mean()# 删除高缺失率字段 high_missing_cols = missing_pct[missing_pct > self.config['missing_threshold']].indexiflen(high_missing_cols) > 0: df_clean = df_clean.drop(columns=high_missing_cols)self.log_operation("缺失值处理", f"删除缺失率>{self.config['missing_threshold']*100}%的字段: {list(high_missing_cols)}")# 填充其他缺失值for col in df_clean.columns:if df_clean[col].isnull().any():# 数值型字段用中位数填充if pd.api.types.is_numeric_dtype(df_clean[col]): fill_value = df_clean[col].median() df_clean[col] = df_clean[col].fillna(fill_value)self.log_operation("缺失值处理", f"字段 {col}: 数值型,用中位数 {fill_value:.2f} 填充")# 分类型字段用众数填充else: fill_value = df_clean[col].mode()[0] ifnot df_clean[col].mode().empty else'Unknown' df_clean[col] = df_clean[col].fillna(fill_value)self.log_operation("缺失值处理", f"字段 {col}: 分类型,用众数 '{fill_value}' 填充")elif strategy == 'drop':# 删除包含缺失值的行(仅当缺失值较少时使用) initial_missing = df_clean.isnull().sum().sum() df_clean = df_clean.dropna() rows_dropped = initial_shape[0] - df_clean.shape[0]self.log_operation("缺失值处理", f"删除包含缺失值的行,删除 {rows_dropped} 行")elif strategy == 'fill':# 自定义填充策略 fill_strategies = {'age': df_clean['age'].median(),'income': df_clean['income'].median(),'education': 'Unknown','marital_status': 'Unknown' }for col, fill_value in fill_strategies.items():if col in df_clean.columns and df_clean[col].isnull().any(): df_clean[col] = df_clean[col].fillna(fill_value)self.log_operation("缺失值处理", f"字段 {col}: 用 '{fill_value}' 填充") final_shape = df_clean.shapeself.log_operation("缺失值处理", f"处理完成,形状从 {initial_shape} 变为 {final_shape}")return df_cleandefhandle_outliers(self, df, numeric_cols=None):""" 处理异常值 Parameters: ----------- df : pd.DataFrame 原始数据 numeric_cols : list, 可选 需要处理的数值型字段列表,如为None则自动检测 Returns: -------- pd.DataFrame 处理后的数据 """self.log_operation("异常值处理", "开始处理异常值") df_clean = df.copy()if numeric_cols isNone:# 自动检测数值型字段 numeric_cols = df_clean.select_dtypes(include=[np.number]).columns.tolist()self.log_operation("异常值处理", f"处理的数值型字段: {numeric_cols}") outliers_summary = {}for col in numeric_cols:if col notin df_clean.columns:continue initial_count = len(df_clean) initial_mean = df_clean[col].mean() initial_std = df_clean[col].std()ifself.config['outlier_method'] == 'IQR':# IQR方法 Q1 = df_clean[col].quantile(0.25) Q3 = df_clean[col].quantile(0.75) IQR = Q3 - Q1 lower_bound = Q1 - 1.5 * IQR upper_bound = Q3 + 1.5 * IQR# 识别异常值 outliers = df_clean[(df_clean[col] < lower_bound) | (df_clean[col] > upper_bound)] outlier_count = len(outliers)# 截断处理(用边界值替换) df_clean[col] = df_clean[col].clip(lower=lower_bound, upper=upper_bound) outliers_summary[col] = {'method': 'IQR','lower_bound': lower_bound,'upper_bound': upper_bound,'outlier_count': outlier_count,'outlier_pct': round(outlier_count / initial_count * 100, 2) }elifself.config['outlier_method'] == '3sigma':# 3σ方法 mean = df_clean[col].mean() std = df_clean[col].std() lower_bound = mean - 3 * std upper_bound = mean + 3 * std# 识别异常值 outliers = df_clean[(df_clean[col] < lower_bound) | (df_clean[col] > upper_bound)] outlier_count = len(outliers)# 截断处理 df_clean[col] = df_clean[col].clip(lower=lower_bound, upper=upper_bound) outliers_summary[col] = {'method': '3sigma','lower_bound': lower_bound,'upper_bound': upper_bound,'outlier_count': outlier_count,'outlier_pct': round(outlier_count / initial_count * 100, 2) } final_mean = df_clean[col].mean() final_std = df_clean[col].std()self.log_operation("异常值处理", f"字段 {col}: 发现 {outlier_count} 个异常值 ({outliers_summary[col]['outlier_pct']}%),已处理")self.log_operation("异常值处理", f"字段 {col}: 均值从 {initial_mean:.2f} 变为 {final_mean:.2f},标准差从 {initial_std:.2f} 变为 {final_std:.2f}")return df_clean, outliers_summarydefvalidate_basic_rules(self, df):""" 验证基本业务规则 Parameters: ----------- df : pd.DataFrame 原始数据 Returns: -------- pd.DataFrame 验证后的数据 dict 验证结果统计 """self.log_operation("业务规则验证", "开始验证基本业务规则") df_clean = df.copy() validation_results = {'规则违反记录': {},'修正记录': {},'删除记录': {} }# 规则1: 年龄范围验证if'age'in df_clean.columns: invalid_age = df_clean[(df_clean['age'] < self.config['min_age']) | (df_clean['age'] > self.config['max_age'])]iflen(invalid_age) > 0: validation_results['规则违反记录']['年龄范围'] = {'违反数量': len(invalid_age),'描述': f"年龄不在{self.config['min_age']}-{self.config['max_age']}岁之间",'样本ID': invalid_age.index.tolist()[:10] # 只显示前10个 }# 标记为无效,但不立即删除(可在后续处理) df_clean['age_valid'] = ~df_clean.index.isin(invalid_age.index)self.log_operation("业务规则验证", f"年龄规则违反: {len(invalid_age)} 个样本年龄不在有效范围内")# 规则2: 收入范围验证if'income'in df_clean.columns: invalid_income = df_clean[(df_clean['income'] < self.config['min_income']) | (df_clean['income'] > self.config['max_income'])]iflen(invalid_income) > 0: validation_results['规则违反记录']['收入范围'] = {'违反数量': len(invalid_income),'描述': f"收入不在{self.config['min_income']}-{self.config['max_income']}之间",'样本ID': invalid_income.index.tolist()[:10] }# 对于异常高收入,可以用上限值截断 high_income = df_clean[df_clean['income'] > self.config['max_income']]iflen(high_income) > 0: df_clean.loc[high_income.index, 'income'] = self.config['max_income'] validation_results['修正记录']['高收入修正'] = {'修正数量': len(high_income),'修正方法': f'用上限值{self.config["max_income"]}替换' }self.log_operation("业务规则验证", f"收入修正: {len(high_income)} 个样本收入过高,已用上限值替换")# 规则3: 工作年限不应大于年龄-18ifall(col in df_clean.columns for col in ['age', 'work_years']): invalid_work = df_clean[df_clean['work_years'] > (df_clean['age'] - 18)]iflen(invalid_work) > 0: validation_results['规则违反记录']['工作年限'] = {'违反数量': len(invalid_work),'描述': "工作年限大于(年龄-18)",'样本ID': invalid_work.index.tolist()[:10] }# 修正:重新计算合理的工作年限 df_clean.loc[invalid_work.index, 'work_years'] = df_clean.loc[invalid_work.index, 'age'] - 18 validation_results['修正记录']['工作年限修正'] = {'修正数量': len(invalid_work),'修正方法': '重新计算为(年龄-18)' }self.log_operation("业务规则验证", f"工作年限修正: {len(invalid_work)} 个样本工作年限不合理,已重新计算")# 规则4: 手机号格式验证if'mobile'in df_clean.columns:# 简单手机号验证:11位数字,以1开头 mobile_pattern = r'^1[3-9]\d{9}$' invalid_mobile = df_clean[~df_clean['mobile'].astype(str).str.match(mobile_pattern, na=False)]iflen(invalid_mobile) > 0: validation_results['规则违反记录']['手机号格式'] = {'违反数量': len(invalid_mobile),'描述': "手机号格式不正确",'样本ID': invalid_mobile.index.tolist()[:10] }self.log_operation("业务规则验证", f"手机号格式违反: {len(invalid_mobile)} 个样本手机号格式不正确")# 规则5: 身份证号验证if'id_card'in df_clean.columns:# 简单身份证验证:18位,最后一位可能是X id_pattern = r'^\d{17}[\dXx]$' invalid_id = df_clean[~df_clean['id_card'].astype(str).str.match(id_pattern, na=False)]iflen(invalid_id) > 0: validation_results['规则违反记录']['身份证号格式'] = {'违反数量': len(invalid_id),'描述': "身份证号格式不正确",'样本ID': invalid_id.index.tolist()[:10] }# 严重错误:身份证号格式错误通常意味着数据质量问题# 可以考虑删除这些记录iflen(invalid_id) / len(df_clean) < 0.01: # 如果错误率<1%,可以删除 df_clean = df_clean.drop(invalid_id.index) validation_results['删除记录']['身份证号错误'] = {'删除数量': len(invalid_id),'删除原因': "身份证号格式错误" }self.log_operation("业务规则验证", f"删除记录: {len(invalid_id)} 个样本身份证号格式错误,已删除")self.log_operation("业务规则验证", "业务规则验证完成")return df_clean, validation_resultsdeffeature_engineering(self, df):""" 基础特征工程 Parameters: ----------- df : pd.DataFrame 原始数据 Returns: -------- pd.DataFrame 特征工程后的数据 """self.log_operation("特征工程", "开始基础特征工程") df_fe = df.copy() new_features = []# 1. 创建比率特征ifall(col in df_fe.columns for col in ['total_debt', 'annual_income']): df_fe['debt_to_income_ratio'] = df_fe['total_debt'] / df_fe['annual_income'].replace(0, 1) new_features.append('debt_to_income_ratio')self.log_operation("特征工程", "创建特征: debt_to_income_ratio")ifall(col in df_fe.columns for col in ['monthly_payment', 'monthly_income']): df_fe['payment_to_income_ratio'] = df_fe['monthly_payment'] / df_fe['monthly_income'].replace(0, 1) new_features.append('payment_to_income_ratio')self.log_operation("特征工程", "创建特征: payment_to_income_ratio")# 2. 创建稳定性特征ifall(col in df_fe.columns for col in ['current_job_years', 'age']): df_fe['job_stability'] = df_fe['current_job_years'] / (df_fe['age'] - 18).replace(0, 1) new_features.append('job_stability')self.log_operation("特征工程", "创建特征: job_stability")if'address_years'in df_fe.columns: df_fe['address_stability'] = df_fe['address_years'] / 10# 按10年标准化 new_features.append('address_stability')self.log_operation("特征工程", "创建特征: address_stability")# 3. 创建时间特征if'application_date'in df_fe.columns:try: df_fe['application_date'] = pd.to_datetime(df_fe['application_date']) df_fe['application_month'] = df_fe['application_date'].dt.month df_fe['application_dayofweek'] = df_fe['application_date'].dt.dayofweek df_fe['application_hour'] = df_fe['application_date'].dt.hour new_features.extend(['application_month', 'application_dayofweek', 'application_hour'])self.log_operation("特征工程", "创建时间特征")except:self.log_operation("特征工程", "警告: 日期格式转换失败")# 4. 创建交叉特征ifall(col in df_fe.columns for col in ['age', 'income_level']):# 假设income_level是分类变量 df_fe['age_income_interaction'] = df_fe['age'] * df_fe['income_level'].astype('category').cat.codes new_features.append('age_income_interaction')self.log_operation("特征工程", "创建特征: age_income_interaction")# 5. 创建风险分类特征if'credit_score'in df_fe.columns:# 根据信用评分分段 bins = [0, 550, 650, 750, 850] labels = ['高风险', '中高风险', '中等风险', '低风险'] df_fe['risk_level'] = pd.cut(df_fe['credit_score'], bins=bins, labels=labels) new_features.append('risk_level')self.log_operation("特征工程", "创建特征: risk_level")self.log_operation("特征工程", f"共创建 {len(new_features)} 个新特征: {new_features}")return df_fedefencode_categorical_features(self, df, cat_cols=None):""" 编码分类特征 Parameters: ----------- df : pd.DataFrame 原始数据 cat_cols : list, 可选 分类特征列表,如为None则自动检测 Returns: -------- pd.DataFrame 编码后的数据 dict 编码映射 """self.log_operation("特征编码", "开始编码分类特征") df_encoded = df.copy() encoding_maps = {}if cat_cols isNone:# 自动检测分类特征 cat_cols = df_encoded.select_dtypes(include=['object', 'category']).columns.tolist()# 过滤掉可能不是真正分类的字段(如ID) exclude_cols = ['id', 'customer_id', 'application_id', 'mobile', 'id_card'] cat_cols = [col for col in cat_cols if col notin exclude_cols]self.log_operation("特征编码", f"需要编码的分类特征: {cat_cols}")for col in cat_cols:if col notin df_encoded.columns:continue# 处理缺失值if df_encoded[col].isnull().any(): df_encoded[col] = df_encoded[col].fillna('Missing')self.log_operation("特征编码", f"字段 {col}: 填充缺失值为 'Missing'")# 对于基数较低的分类变量,使用标签编码 unique_count = df_encoded[col].nunique()if unique_count <= 10:# 标签编码from sklearn.preprocessing import LabelEncoder le = LabelEncoder() df_encoded[col + '_encoded'] = le.fit_transform(df_encoded[col]) encoding_maps[col] = {'method': 'LabelEncoding','mapping': dict(zip(le.classes_, le.transform(le.classes_))) }self.log_operation("特征编码", f"字段 {col}: 标签编码,{unique_count} 个类别")else:# 对于基数高的分类变量,使用频率编码或目标编码(需要目标变量)# 这里使用频率编码作为示例 freq = df_encoded[col].value_counts(normalize=True) df_encoded[col + '_freq'] = df_encoded[col].map(freq) encoding_maps[col] = {'method': 'FrequencyEncoding','mapping': freq.to_dict() }self.log_operation("特征编码", f"字段 {col}: 频率编码,{unique_count} 个类别")self.log_operation("特征编码", "分类特征编码完成")return df_encoded, encoding_mapsdefdetect_duplicates(self, df, id_cols=None):""" 检测重复数据 Parameters: ----------- df : pd.DataFrame 原始数据 id_cols : list, 可选 用于识别重复的字段列表 Returns: -------- pd.DataFrame 去重后的数据 dict 重复检测结果 """self.log_operation("重复检测", "开始检测重复数据")if id_cols isNone:# 默认使用可能识别唯一客户的字段 id_cols = ['id_card', 'mobile', 'customer_id'] id_cols = [col for col in id_cols if col in df.columns]ifnot id_cols:self.log_operation("重复检测", "警告: 没有可用的ID字段,跳过重复检测")return df, {} duplicate_results = {}# 方法1: 基于关键字段的精确匹配for col in id_cols: duplicates = df[df.duplicated(subset=[col], keep='first')]iflen(duplicates) > 0: duplicate_results[f'{col}_duplicates'] = {'重复数量': len(duplicates),'重复值示例': duplicates[col].unique()[:5].tolist() }self.log_operation("重复检测", f"字段 {col}: 发现 {len(duplicates)} 个重复值")# 方法2: 基于多个字段的组合重复iflen(id_cols) >= 2: combo_duplicates = df[df.duplicated(subset=id_cols, keep='first')]iflen(combo_duplicates) > 0: duplicate_results['combo_duplicates'] = {'重复数量': len(combo_duplicates),'重复字段': id_cols }self.log_operation("重复检测", f"组合字段 {id_cols}: 发现 {len(combo_duplicates)} 个重复记录")# 删除重复记录(保留第一个) initial_count = len(df) df_deduped = df.drop_duplicates(subset=id_cols, keep='first') removed_count = initial_count - len(df_deduped)if removed_count > 0: duplicate_results['removal_summary'] = {'删除数量': removed_count,'保留数量': len(df_deduped) }self.log_operation("重复检测", f"删除 {removed_count} 个重复记录,保留 {len(df_deduped)} 个唯一记录")return df_deduped, duplicate_resultsdefsave_cleaned_data(self, df, file_path, file_type='csv', **kwargs):""" 保存清洗后的数据 Parameters: ----------- df : pd.DataFrame 清洗后的数据 file_path : str 保存路径 file_type : str 文件类型:'csv', 'excel', 'parquet' **kwargs : dict pandas保存参数 """self.log_operation("数据保存", f"保存清洗后的数据到 {file_path}")try:if file_type == 'csv': df.to_csv(file_path, index=False, **kwargs)elif file_type == 'excel': df.to_excel(file_path, index=False, **kwargs)elif file_type == 'parquet': df.to_parquet(file_path, index=False, **kwargs)else:raise ValueError(f"不支持的文件类型: {file_type}")self.log_operation("数据保存", "数据保存成功")except Exception as e:self.log_operation("数据保存", f"保存失败: {str(e)}")raisedefgenerate_summary_report(self):""" 生成清洗总结报告 Returns: -------- dict 清洗总结报告 """ report = {'清洗时间': datetime.now().strftime('%Y-%m-%d %H:%M:%S'),'清洗步骤': len(self.cleaning_log),'清洗日志摘要': self.cleaning_log[-10:], # 最后10条日志'配置参数': self.config }return report# ============================================================================# 使用示例# ============================================================================defexample_usage():"""使用示例"""# 1. 初始化清洗器 config = {'missing_threshold': 0.3, # 缺失率超过30%的字段将被删除'outlier_method': 'IQR','min_age': 20,'max_age': 65,'save_log': True } cleaner = RiskDataCleaner(config)# 2. 加载数据(示例数据)# 创建示例数据 np.random.seed(42) n_samples = 1000 example_data = pd.DataFrame({'customer_id': range(1, n_samples + 1),'age': np.random.randint(18, 70, n_samples),'income': np.random.normal(15000, 5000, n_samples).clip(3000, 50000),'education': np.random.choice(['高中', '大专', '本科', '硕士', '博士', None], n_samples, p=[0.2, 0.3, 0.3, 0.15, 0.04, 0.01]),'work_years': np.random.randint(0, 40, n_samples),'total_debt': np.random.exponential(50000, n_samples).clip(0, 300000),'annual_income': np.random.normal(200000, 50000, n_samples).clip(50000, 500000),'credit_score': np.random.randint(300, 850, n_samples),'mobile': ['138' + str(np.random.randint(10000000, 99999999)) for _ inrange(n_samples)],'application_date': pd.date_range('2023-01-01', periods=n_samples, freq='H') })# 添加一些缺失值和异常值 example_data.loc[np.random.choice(n_samples, 50), 'age'] = np.nan example_data.loc[np.random.choice(n_samples, 30), 'income'] = 1000000# 异常高收入 example_data.loc[100, 'work_years'] = 60# 不合理的工作年限# 保存示例数据 example_data.to_csv('example_risk_data.csv', index=False)print("=" * 60)print("金融风控数据清洗示例")print("=" * 60)# 3. 加载数据 df = cleaner.load_data('example_risk_data.csv')# 4. 分析数据质量 quality_report = cleaner.analyze_data_quality(df)print("\n数据质量报告:")print(f"数据形状: {quality_report['数据形状']}")print(f"总缺失值: {quality_report['缺失值统计']['总缺失值']}")# 5. 处理缺失值 df_clean = cleaner.handle_missing_values(df, strategy='auto')# 6. 处理异常值 df_clean, outliers_summary = cleaner.handle_outliers(df_clean)# 7. 验证业务规则 df_clean, validation_results = cleaner.validate_basic_rules(df_clean)# 8. 特征工程 df_fe = cleaner.feature_engineering(df_clean)# 9. 编码分类特征 df_encoded, encoding_maps = cleaner.encode_categorical_features(df_fe)# 10. 检测重复数据 df_final, duplicate_results = cleaner.detect_duplicates(df_encoded, id_cols=['customer_id'])# 11. 保存清洗后的数据 cleaner.save_cleaned_data(df_final, 'cleaned_risk_data.csv')# 12. 保存清洗日志 cleaner.save_log()# 13. 生成总结报告 final_report = cleaner.generate_summary_report()print("\n" + "=" * 60)print("清洗完成!")print(f"原始数据形状: {df.shape}")print(f"清洗后数据形状: {df_final.shape}")print(f"数据清洗日志已保存到: {cleaner.config['log_file']}")print(f"清洗后数据已保存到: cleaned_risk_data.csv")print("=" * 60)return df_final, final_report# ============================================================================# 快速清洗函数(一键式清洗)# ============================================================================defquick_clean_data(file_path, target_file=None, config=None):""" 快速数据清洗(一键式) Parameters: ----------- file_path : str 原始数据文件路径 target_file : str, 可选 清洗后数据保存路径 config : dict, 可选 清洗配置 Returns: -------- pd.DataFrame 清洗后的数据 """if config isNone: config = {'missing_threshold': 0.5,'outlier_method': 'IQR','min_age': 18,'max_age': 70,'save_log': True,'log_file': 'quick_clean_log.txt' }# 初始化清洗器 cleaner = RiskDataCleaner(config)# 确定文件类型 file_type = 'csv'if file_path.endswith('.xlsx') or file_path.endswith('.xls'): file_type = 'excel'elif file_path.endswith('.parquet'): file_type = 'parquet'# 加载数据 df = cleaner.load_data(file_path, file_type=file_type)# 数据质量分析 quality_report = cleaner.analyze_data_quality(df)# 完整清洗流程 df_clean = cleaner.handle_missing_values(df, strategy='auto') df_clean, _ = cleaner.handle_outliers(df_clean) df_clean, _ = cleaner.validate_basic_rules(df_clean) df_clean = cleaner.feature_engineering(df_clean) df_clean, _ = cleaner.encode_categorical_features(df_clean) df_clean, _ = cleaner.detect_duplicates(df_clean)# 保存数据if target_file isNone: target_file = file_path.replace('.csv', '_cleaned.csv').replace('.xlsx', '_cleaned.xlsx') cleaner.save_cleaned_data(df_clean, target_file) cleaner.save_log()print(f"快速清洗完成!")print(f"原始数据: {file_path}")print(f"清洗后数据: {target_file}")print(f"清洗日志: {cleaner.config['log_file']}")return df_clean# ============================================================================# 主程序入口# ============================================================================if __name__ == "__main__":# 运行示例 cleaned_data, report = example_usage()# 或者使用快速清洗# cleaned_data = quick_clean_data('your_data.csv')二、核心清洗函数独立版本
如果您只需要特定的清洗功能,这里提供独立版本:
# 1. 缺失值处理独立函数defhandle_missing_values_simple(df, numeric_strategy='median', categorical_strategy='mode'):""" 简易缺失值处理 Parameters: ----------- df : pd.DataFrame 原始数据 numeric_strategy : str 数值型字段填充策略:'mean', 'median', 'zero' categorical_strategy : str 分类型字段填充策略:'mode', 'unknown' Returns: -------- pd.DataFrame 处理后的数据 """ df_clean = df.copy()for col in df_clean.columns:if df_clean[col].isnull().any():# 数值型字段if pd.api.types.is_numeric_dtype(df_clean[col]):if numeric_strategy == 'mean': fill_value = df_clean[col].mean()elif numeric_strategy == 'median': fill_value = df_clean[col].median()elif numeric_strategy == 'zero': fill_value = 0else: fill_value = df_clean[col].median() df_clean[col] = df_clean[col].fillna(fill_value)print(f"字段 {col}: 数值型,用 {numeric_strategy}({fill_value:.2f}) 填充")# 分类型字段else:if categorical_strategy == 'mode': fill_value = df_clean[col].mode()[0] ifnot df_clean[col].mode().empty else'Unknown'elif categorical_strategy == 'unknown': fill_value = 'Unknown'else: fill_value = 'Unknown' df_clean[col] = df_clean[col].fillna(fill_value)print(f"字段 {col}: 分类型,用 '{fill_value}' 填充")return df_clean# 2. 异常值处理独立函数defdetect_and_treat_outliers_iqr(df, columns=None):""" 使用IQR方法检测和处理异常值 Parameters: ----------- df : pd.DataFrame 原始数据 columns : list, 可选 需要处理的字段列表 Returns: -------- pd.DataFrame 处理后的数据 dict 异常值统计 """ df_clean = df.copy() outlier_stats = {}if columns isNone: columns = df_clean.select_dtypes(include=[np.number]).columnsfor col in columns:if col notin df_clean.columns:continue# 计算IQR Q1 = df_clean[col].quantile(0.25) Q3 = df_clean[col].quantile(0.75) IQR = Q3 - Q1# 定义异常值边界 lower_bound = Q1 - 1.5 * IQR upper_bound = Q3 + 1.5 * IQR# 检测异常值 outliers = df_clean[(df_clean[col] < lower_bound) | (df_clean[col] > upper_bound)] outlier_count = len(outliers)# 用边界值截断(Winsorization) df_clean[col] = df_clean[col].clip(lower=lower_bound, upper=upper_bound) outlier_stats[col] = {'Q1': Q1,'Q3': Q3,'IQR': IQR,'lower_bound': lower_bound,'upper_bound': upper_bound,'outlier_count': outlier_count,'outlier_pct': round(outlier_count / len(df_clean) * 100, 2) }if outlier_count > 0:print(f"字段 {col}: 发现 {outlier_count} 个异常值({outlier_stats[col]['outlier_pct']}%),已处理")return df_clean, outlier_stats# 3. 金融风控专用业务规则验证defvalidate_financial_rules(df):""" 金融风控专用业务规则验证 Parameters: ----------- df : pd.DataFrame 原始数据 Returns: -------- pd.DataFrame 验证后的数据 dict 验证结果 """ validation_results = {'issues_found': 0,'issues_fixed': 0,'details': [] } df_clean = df.copy()# 规则1: 年龄合理性if'age'in df_clean.columns: invalid_age = df_clean[(df_clean['age'] < 18) | (df_clean['age'] > 70)]iflen(invalid_age) > 0: validation_results['issues_found'] += len(invalid_age) validation_results['details'].append({'field': 'age','issue': f"{len(invalid_age)}个样本年龄不在18-70岁之间",'action': '标记为异常' }) df_clean['age_flag'] = df_clean['age'].apply(lambda x: 1if (x < 18or x > 70) else0)# 规则2: 收入负债比ifall(col in df_clean.columns for col in ['monthly_debt', 'monthly_income']): df_clean['dti_ratio'] = df_clean['monthly_debt'] / df_clean['monthly_income'].replace(0, 1) high_dti = df_clean[df_clean['dti_ratio'] > 0.5] # 负债收入比超过50%iflen(high_dti) > 0: validation_results['issues_found'] += len(high_dti) validation_results['details'].append({'field': 'dti_ratio','issue': f"{len(high_dti)}个样本负债收入比>50%",'action': '标记为高风险' }) df_clean['high_dti_flag'] = (df_clean['dti_ratio'] > 0.5).astype(int)# 规则3: 工作稳定性if'job_tenure'in df_clean.columns: short_tenure = df_clean[df_clean['job_tenure'] < 6] # 工作少于6个月iflen(short_tenure) > 0: validation_results['issues_found'] += len(short_tenure) validation_results['details'].append({'field': 'job_tenure','issue': f"{len(short_tenure)}个样本工作年限<6个月",'action': '标记为不稳定' }) df_clean['job_stability_flag'] = (df_clean['job_tenure'] < 6).astype(int)# 规则4: 申请频率异常if'application_count_30d'in df_clean.columns: frequent_applicants = df_clean[df_clean['application_count_30d'] > 5] # 30天内申请超过5次iflen(frequent_applicants) > 0: validation_results['issues_found'] += len(frequent_applicants) validation_results['issues_fixed'] += len(frequent_applicants) validation_results['details'].append({'field': 'application_count_30d','issue': f"{len(frequent_applicants)}个样本30天内申请超过5次",'action': '标记为可疑申请' }) df_clean['suspicious_app_flag'] = (df_clean['application_count_30d'] > 5).astype(int) validation_results['total_samples'] = len(df_clean)return df_clean, validation_results# 4. 特征工程模板defcreate_risk_features(df):""" 创建风控特征 Parameters: ----------- df : pd.DataFrame 原始数据 Returns: -------- pd.DataFrame 包含新特征的数据 """ df_fe = df.copy()# 基础比率特征ifall(col in df_fe.columns for col in ['total_liabilities', 'annual_income']): df_fe['liability_to_income'] = df_fe['total_liabilities'] / df_fe['annual_income'].replace(0, 1)ifall(col in df_fe.columns for col in ['credit_card_utilization', 'credit_limit']): df_fe['utilization_rate'] = df_fe['credit_card_utilization'] / df_fe['credit_limit'].replace(0, 1)# 稳定性特征if'current_address_months'in df_fe.columns: df_fe['address_stability'] = df_fe['current_address_months'] / 12# 转换为年if'current_employer_months'in df_fe.columns: df_fe['employment_stability'] = df_fe['current_employer_months'] / 12# 行为特征if'num_credit_inquiries_6m'in df_fe.columns: df_fe['inquiry_intensity'] = df_fe['num_credit_inquiries_6m'] / 6# 月均查询次数# 时间特征if'application_datetime'in df_fe.columns: df_fe['application_hour'] = pd.to_datetime(df_fe['application_datetime']).dt.hour df_fe['is_weekend'] = pd.to_datetime(df_fe['application_datetime']).dt.dayofweek >= 5 df_fe['is_night'] = df_fe['application_hour'].between(0, 5) | df_fe['application_hour'].between(22, 23)# 组合特征ifall(col in df_fe.columns for col in ['age', 'income']): df_fe['age_income_combo'] = df_fe['age'] * (df_fe['income'] / 10000)return df_fe三、实战使用案例
"""实战案例:消费金融公司数据清洗"""import pandas as pdimport numpy as np# 案例1:消费贷申请数据清洗defclean_consumer_loan_data(data_path):""" 清洗消费贷申请数据 """# 读取数据 df = pd.read_csv(data_path)print(f"原始数据形状: {df.shape}")print(f"字段列表: {list(df.columns)}")# 1. 处理缺失值from RiskDataCleaner import handle_missing_values_simple df = handle_missing_values_simple(df, numeric_strategy='median', categorical_strategy='mode')# 2. 处理异常值 numeric_cols = ['age', 'monthly_income', 'loan_amount', 'work_years'] df, outlier_stats = detect_and_treat_outliers_iqr(df, numeric_cols)# 3. 验证业务规则 df, validation_results = validate_financial_rules(df)# 4. 创建风控特征 df = create_risk_features(df)# 5. 保存清洗后的数据 df.to_csv('cleaned_consumer_loan_data.csv', index=False)print(f"清洗完成!")print(f"清洗后数据形状: {df.shape}")print(f"发现的问题数量: {validation_results['issues_found']}")return df# 案例2:信用卡申请数据清洗defclean_credit_card_data(df):""" 清洗信用卡申请数据 """# 信用卡数据特定规则 rules = {'min_age': 18,'max_age': 65,'min_income': 2000, # 月收入最低要求'max_credit_limit': 200000, # 最高信用额度'min_work_years': 0.5# 最低工作年限(年) }# 应用规则 mask_valid = ( (df['age'] >= rules['min_age']) & (df['age'] <= rules['max_age']) & (df['monthly_income'] >= rules['min_income']) & (df['work_years'] >= rules['min_work_years']) ) df_valid = df[mask_valid].copy() df_invalid = df[~mask_valid].copy()print(f"有效申请: {len(df_valid)}")print(f"无效申请: {len(df_invalid)}")# 对有效申请进一步处理iflen(df_valid) > 0:# 计算信用评分(简化版) df_valid['credit_score_simple'] = ( (df_valid['monthly_income'] / 1000 * 10) + (df_valid['work_years'] * 5) + (df_valid['age'] * 0.5) ).astype(int).clip(300, 850)# 根据评分分配信用额度defassign_credit_limit(score, income):if score >= 750:returnmin(income * 24, rules['max_credit_limit'])elif score >= 650:returnmin(income * 18, rules['max_credit_limit'])elif score >= 550:returnmin(income * 12, rules['max_credit_limit'])else:returnmin(income * 6, rules['max_credit_limit']) df_valid['suggested_limit'] = df_valid.apply(lambda row: assign_credit_limit(row['credit_score_simple'], row['monthly_income']), axis=1 )return df_valid, df_invalid# 案例3:批量清洗多个文件defbatch_clean_data(file_list, output_dir='cleaned_data'):""" 批量清洗多个数据文件 """import osifnot os.path.exists(output_dir): os.makedirs(output_dir) results = []for file_path in file_list:print(f"\n处理文件: {file_path}")try:# 使用快速清洗 df_clean = quick_clean_data(file_path)# 保存结果 filename = os.path.basename(file_path) output_path = os.path.join(output_dir, f"cleaned_{filename}") df_clean.to_csv(output_path, index=False) results.append({'file': filename,'status': 'success','original_rows': 'N/A', # 需要从原始数据获取'cleaned_rows': len(df_clean),'output_path': output_path })print(f"成功清洗并保存到: {output_path}")except Exception as e:print(f"处理失败: {str(e)}") results.append({'file': os.path.basename(file_path),'status': 'failed','error': str(e) })# 生成批量处理报告 report_df = pd.DataFrame(results) report_path = os.path.join(output_dir, 'batch_clean_report.csv') report_df.to_csv(report_path, index=False)print(f"\n批量处理完成!")print(f"处理报告已保存到: {report_path}")return results四、模板使用建议
1. 快速上手步骤
# 第一步:导入模板from RiskDataCleaner import RiskDataCleaner# 第二步:初始化清洗器cleaner = RiskDataCleaner()# 第三步:加载数据df = cleaner.load_data('your_data.csv')# 第四步:一键清洗(使用完整流程)df_clean = cleaner.full_clean_pipeline(df)# 第五步:保存结果cleaner.save_cleaned_data(df_clean, 'cleaned_data.csv')2. 自定义配置
# 根据业务需求自定义配置custom_config = {'missing_threshold': 0.3, # 更严格的缺失值处理'outlier_method': '3sigma', # 使用3σ方法'min_age': 20, # 提高最低年龄要求'max_age': 60, # 降低最高年龄要求'min_income': 3000, # 设置最低收入要求'save_log': True,'log_file': 'my_cleaning_log.txt'}cleaner = RiskDataCleaner(custom_config)3. 处理特定问题
# 只处理缺失值df_no_missing = cleaner.handle_missing_values(df, strategy='auto')# 只处理异常值df_no_outliers, stats = cleaner.handle_outliers(df)# 只做特征工程df_with_features = cleaner.feature_engineering(df)五、最佳实践建议
1. 数据备份:清洗前始终备份原始数据 2. 逐步验证:每步清洗后检查数据质量 3. 业务理解:清洗规则要基于业务逻辑 4. 文档记录:记录所有清洗决策和原因 5. 版本控制:保存不同版本的清洗脚本和数据
此模板覆盖了金融风控数据清洗的85%常见场景,可根据具体业务需求进行调整和扩展。
欢迎添加:
夜雨聆风