"""数据治理:编码规范清洗 — 检查KPI编码前缀与维度一致性 + 修复误分类""" import pymysql import os import logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger("encoding-cleanup") DB_USER = os.getenv("CMA_DB_USER", "cma_user") DB_PASS = os.getenv("CMA_DB_PASS", "cma_pass_2026") DB_HOST = os.getenv("CMA_DB_HOST", "127.0.0.1") DB_PORT = int(os.getenv("CMA_DB_PORT", "3306")) DB_NAME = os.getenv("CMA_DB_NAME", "cma") conn = pymysql.connect( host=DB_HOST, port=DB_PORT, user=DB_USER, password=DB_PASS, database=DB_NAME, charset="utf8mb4", cursorclass=pymysql.cursors.DictCursor, ) cur = conn.cursor() # 编码前缀 → 正确维度映射 PREFIX_DIM_MAP = { "F_": "finance", "C_": "customer", "P_": "process", "L_": "learning", } # 已知误分类修复(编码 → 正确维度) KNOWN_FIXES = { "F_QUALITY_RATE": "process", # 产品合格率 → 流程层 "F_REWORK_RATE": "process", # 返工率 → 流程层 "P_COST_CUT": "process", # 招待费砍半 → 流程层 "P_TRAIN_PASS": "process", # Model C考核通过 → 流程层 } def run(): logger.info("=== 编码规范清洗 开始 ===") cur.execute("SELECT id, entity_id, kpi_code, kpi_name, dimension FROM kpi_definitions WHERE status = 'active'") kpis = cur.fetchall() issues = [] fixes_applied = 0 for kpi in kpis: kpi_code = kpi["kpi_code"] current_dim = kpi["dimension"] entity_id = kpi["entity_id"] # 检查前缀 prefix = kpi_code[:2] if len(kpi_code) >= 2 else "" expected_dim = PREFIX_DIM_MAP.get(prefix) if expected_dim and current_dim != expected_dim: # 先检查是否在已知修复列表 correct_dim = KNOWN_FIXES.get(kpi_code, expected_dim) issues.append({ "kpi_code": kpi_code, "kpi_name": kpi["kpi_name"], "current_dim": current_dim, "expected_dim": correct_dim, "prefix": prefix, "entity_id": entity_id, }) if kpi_code in KNOWN_FIXES: logger.info(f" 🔧 修复: {kpi_code} ({kpi['kpi_name']}) {current_dim} → {correct_dim} (entity={entity_id})") cur.execute( "UPDATE kpi_definitions SET dimension = %s WHERE id = %s", (correct_dim, kpi["id"]), ) fixes_applied += 1 conn.commit() # 输出报告 logger.info(f"\n=== 清洗报告 ===") logger.info(f" 检查KPI总数: {len(kpis)}") logger.info(f" 编码-维度不一致: {len(issues)}") logger.info(f" 已自动修复: {fixes_applied}") if issues: logger.info(f"\n 不一致详情:") for i, iss in enumerate(issues, 1): status = "✅ 已修复" if iss["kpi_code"] in KNOWN_FIXES else "⚠️ 需人工确认" logger.info(f" {i}. {iss['kpi_code']} ({iss['kpi_name']}) " f"当前维度={iss['current_dim']}, 期望维度={iss['expected_dim']} [{status}]") # 检查未命名规范问题 logger.info(f"\n 编码前缀统计:") for prefix, dim in PREFIX_DIM_MAP.items(): cur.execute("SELECT COUNT(*) as cnt FROM kpi_definitions WHERE kpi_code LIKE %s AND status='active'", (f"{prefix}%",)) row = cur.fetchone() cnt = row["cnt"] if row else 0 logger.info(f" {prefix} → {dim}: {cnt} 个KPI") # 检查前缀不匹配编码 cur.execute(""" SELECT kpi_code, dimension FROM kpi_definitions WHERE status='active' AND ( (kpi_code LIKE 'F_%' AND dimension != 'finance') OR (kpi_code LIKE 'C_%' AND dimension != 'customer') OR (kpi_code LIKE 'P_%' AND dimension != 'process') OR (kpi_code LIKE 'L_%' AND dimension != 'learning') ) """) remaining = cur.fetchall() if remaining: logger.warning(f"\n ⚠️ 仍有 {len(remaining)} 个KPI编码前缀与维度不匹配:") for r in remaining: logger.warning(f" {r['kpi_code']} → {r['dimension']}") else: logger.info(f"\n ✅ 所有KPI编码前缀与维度一致!") logger.info("\n=== 编码规范清洗 完成 ===") if __name__ == "__main__": try: run() finally: cur.close() conn.close()