feat: 数据治理 — 入库约束+元数据卡片+编码清洗+审计看板
This commit is contained in:
@@ -0,0 +1,126 @@
|
||||
"""数据治理:编码规范清洗 — 检查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()
|
||||
@@ -0,0 +1,105 @@
|
||||
"""数据治理 migration: 入库必检约束 + 元数据字段补充"""
|
||||
import pymysql
|
||||
import os
|
||||
import logging
|
||||
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
logger = logging.getLogger("data-governance")
|
||||
|
||||
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()
|
||||
|
||||
|
||||
def run():
|
||||
logger.info("=== 数据治理 Migration 开始 ===")
|
||||
|
||||
# ── 0. 添加缺失列 ──
|
||||
logger.info("[步骤0] 检查并补充缺失列...")
|
||||
|
||||
for col_name, col_def in [
|
||||
("data_source", "ALTER TABLE kpi_definitions ADD COLUMN data_source VARCHAR(500) DEFAULT NULL COMMENT '数据来源' AFTER data_source_type"),
|
||||
("data_owner", "ALTER TABLE kpi_definitions ADD COLUMN data_owner VARCHAR(100) DEFAULT NULL COMMENT '数据责任人' AFTER data_source"),
|
||||
]:
|
||||
cur.execute("SHOW COLUMNS FROM kpi_definitions LIKE %s", (col_name,))
|
||||
if not cur.fetchone():
|
||||
logger.info(f" 添加 {col_name} 列...")
|
||||
cur.execute(col_def)
|
||||
conn.commit()
|
||||
logger.info(f" {col_name} 列已添加")
|
||||
else:
|
||||
logger.info(f" {col_name} 列已存在")
|
||||
|
||||
# ── 1. 修复空值 ──
|
||||
logger.info("[步骤1] 修复空值...")
|
||||
|
||||
for col, default, label in [
|
||||
("target_value", "0", "NULL target_value"),
|
||||
("unit", "'-'", "NULL/empty unit"),
|
||||
("formula", "'待补充'", "NULL/empty formula"),
|
||||
("data_source", "'待补充'", "NULL/empty data_source"),
|
||||
("data_owner", "'待指定'", "NULL/empty data_owner"),
|
||||
]:
|
||||
if col in ("target_value",):
|
||||
r = cur.execute(f"SELECT COUNT(*) as cnt FROM kpi_definitions WHERE {col} IS NULL")
|
||||
else:
|
||||
r = cur.execute(f"SELECT COUNT(*) as cnt FROM kpi_definitions WHERE {col} IS NULL OR {col} = ''")
|
||||
row = cur.fetchone()
|
||||
cnt = row["cnt"] if row else 0
|
||||
logger.info(f" {label}: {cnt} 条")
|
||||
if cnt > 0:
|
||||
if col in ("target_value",):
|
||||
cur.execute(f"UPDATE kpi_definitions SET {col} = {default} WHERE {col} IS NULL")
|
||||
else:
|
||||
cur.execute(f"UPDATE kpi_definitions SET {col} = {default} WHERE {col} IS NULL OR {col} = ''")
|
||||
logger.info(f" 已修复 {cur.rowcount} 条")
|
||||
|
||||
conn.commit()
|
||||
|
||||
# ── 2. 修改列约束为 NOT NULL ──
|
||||
logger.info("[步骤2] 修改列约束...")
|
||||
|
||||
alters = [
|
||||
("target_value", "DECIMAL(15,2) NOT NULL DEFAULT 0"),
|
||||
("unit", "VARCHAR(50) NOT NULL DEFAULT '-'"),
|
||||
("formula", "TEXT NOT NULL"),
|
||||
("data_source", "VARCHAR(500) NOT NULL DEFAULT '待补充'"),
|
||||
("data_owner", "VARCHAR(100) NOT NULL DEFAULT '待指定'"),
|
||||
]
|
||||
for col, col_type in alters:
|
||||
col_comment = {
|
||||
"target_value": "目标值", "unit": "单位", "formula": "计算公式",
|
||||
"data_source": "数据来源", "data_owner": "数据责任人",
|
||||
}[col]
|
||||
try:
|
||||
cur.execute(f"ALTER TABLE kpi_definitions MODIFY {col} {col_type} COMMENT '{col_comment}'")
|
||||
logger.info(f" {col} → {col_type}")
|
||||
except Exception as e:
|
||||
logger.warning(f" {col} 修改失败: {e}")
|
||||
|
||||
conn.commit()
|
||||
|
||||
# ── 3. 验证 ──
|
||||
logger.info("[步骤3] 验证约束...")
|
||||
cur.execute("DESCRIBE kpi_definitions")
|
||||
for col in cur.fetchall():
|
||||
if col['Field'] in ('target_value', 'unit', 'formula', 'data_source', 'data_owner'):
|
||||
logger.info(f" {col['Field']}: Null={col['Null']}, Default={col['Default']}, Type={col['Type']}")
|
||||
|
||||
logger.info("=== 数据治理 Migration 完成 ===")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
run()
|
||||
finally:
|
||||
cur.close()
|
||||
conn.close()
|
||||
Reference in New Issue
Block a user