Files

127 lines
4.3 KiB
Python

"""数据治理:编码规范清洗 — 检查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()