Files
cma-management/backend/app/api/kpis.py
T
Hermes CI Fix 9b973ab9a8 fix: 智能导入——匹配不上的财务报表科目自动创建KPI定义
- 清理科目前缀(一、/减:/加:)后多级匹配
- ⑤仍未匹配→自动创建KPI(PL_001/CF_001/BS_001)
- 避免205条全部跳过的场景
2026-07-29 18:02:50 +08:00

628 lines
23 KiB
Python

"""KPI字典 API"""
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.orm import Session
from sqlalchemy import func
from typing import Optional, List
from datetime import datetime
import json
from app.database import get_db
from app.auth_middleware import require_auth, require_role, filter_kpis_by_role, kpi_visible_dims
from app.models import StrategicMap, MapObjective, KPIDefinition, KPIValue, KPIAlert, OperationLog, Entity, KPICausality, KPIHierarchy
router = APIRouter(prefix="/api/cma/kpis", tags=["KPI字典"],
dependencies=[Depends(require_role("ceo", "finance", "business", "it"))],
)
# 写操作只允许 ceo/finance/it
WRITE_ROLES = Depends(require_role("ceo", "finance", "it"))
@router.get("")
def list_kpis(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
dimension: Optional[str] = None,
keyword: Optional[str] = None,
epic: Optional[str] = None,
category: Optional[str] = None,
entity_id: Optional[int] = None,
db: Session = Depends(get_db),
current_user = Depends(require_auth),
):
query = db.query(KPIDefinition).filter(KPIDefinition.status == "active")
# 角色权限过滤
dims = kpi_visible_dims(current_user.role, db)
if dims:
query = query.filter(KPIDefinition.dimension.in_(dims))
if dimension:
dims_list = [d.strip() for d in dimension.split(',')] if ',' in dimension else [dimension]
query = query.filter(KPIDefinition.dimension.in_(dims_list))
if keyword:
query = query.filter(KPIDefinition.kpi_name.contains(keyword))
if category:
cats_list = [c.strip() for c in category.split(',')] if ',' in category else [category]
query = query.filter(KPIDefinition.category.in_(cats_list))
if entity_id is not None:
query = query.filter(KPIDefinition.entity_id == entity_id)
total = query.count()
kpis = query.order_by(KPIDefinition.kpi_code).offset((page-1)*page_size).limit(page_size).all()
result = {"total": total, "page": page, "page_size": page_size, "data": [kpi_to_dict(k) for k in kpis]}
if entity_id is not None:
ent = db.query(Entity).filter(Entity.id == entity_id).first()
if ent:
result["entity"] = {"id": ent.id, "name": ent.name, "short_name": ent.short_name}
return result
@router.get("/categories")
def get_kpi_categories(current_user = Depends(require_auth), db: Session = Depends(get_db)):
"""获取BSC分类结构(带可见性过滤)"""
from sqlalchemy import func as sa_func
dims = kpi_visible_dims(current_user.role, db)
query = db.query(
KPIDefinition.dimension,
KPIDefinition.category,
sa_func.count(KPIDefinition.id)
).filter(KPIDefinition.status == "active")
if dims:
query = query.filter(KPIDefinition.dimension.in_(dims))
rows = query.group_by(KPIDefinition.dimension, KPIDefinition.category).all()
# 构建树形结构
dim_map = {"finance": "财务", "customer": "客户", "process": "内部流程", "learning": "学习成长"}
cat_map = {
"revenue_growth": "收入增长", "profitability": "盈利水平", "cost_control": "成本费用",
"asset_efficiency": "资产效率", "cash_risk": "现金流风控",
"customer_scale": "客户规模", "customer_concentration": "客户集中度", "customer_satisfaction": "客户满意",
"supply_chain": "供应链效率", "delivery_quality": "交付质量",
"talent_pipeline": "人才梯队", "employee_engagement": "员工敬业", "innovation": "创新改善",
}
tree = []
for dim, cat, cnt in rows:
# 找或创建维度节点
dim_node = next((n for n in tree if n["key"] == dim), None)
if not dim_node:
dim_node = {"key": dim, "label": dim_map.get(dim, dim), "children": []}
tree.append(dim_node)
dim_node["children"].append({
"key": cat,
"label": cat_map.get(cat, cat),
"count": cnt,
})
dim_counts = {}
for d in tree:
dim_counts[d["key"]] = sum(c["count"] for c in d["children"])
d["count"] = dim_counts[d["key"]]
return {"tree": tree, "total": sum(dim_counts.values())}
# ============================================================
# KPI-5: 五档评分引擎(静态路由必须在动态/{kpi_id}之前)
# ============================================================
REVERSE_INDICATORS = ['C_REBATE_RATE', 'P_BUG_RATE', 'P_REWORK_PCT', 'F_DEBT_RATIO',
'P_QUALITY_RATE', 'F_COST_RATIO', 'F_INTEREST_COVER', 'F_QUICK_RATIO']
def _calc_five_tier_score(current_value, target_value, is_reverse=False):
"""五档评分:1-5分(支持正反向指标)"""
if current_value is None or target_value is None or target_value == 0:
return None, "info"
ratio = current_value / target_value
if is_reverse:
# 反向指标:实际值越低越好
if ratio <= 0.5:
return 5, "success" # 远低于目标→卓越
elif ratio <= 0.8:
return 4, "success" # 低于目标→达标
elif ratio <= 1.0:
return 3, "warning" # 接近目标→预警
elif ratio <= 1.2:
return 2, "danger" # 超过目标→危险
else:
return 1, "danger" # 远超目标→失效
else:
if ratio >= 1.2:
return 5, "success" # 卓越
elif ratio >= 1.0:
return 4, "success" # 达标
elif ratio >= 0.8:
return 3, "warning" # 预警
elif ratio >= 0.5:
return 2, "danger" # 危险
else:
return 1, "danger" # 失效
@router.get("/score")
def get_kpi_score(
entity_id: int = Query(1, ge=1),
period: Optional[str] = None,
db: Session = Depends(get_db),
current_user = Depends(require_auth),
):
"""五档评分引擎 - 返回各KPI评分和BSC四层汇总
评分: 5卓越(≥1.2×目标) 4达标(≥目标) 3预警(≥0.8×目标) 2危险(≥0.5×目标) 1失效(<0.5×目标)
"""
# 获取该企业所有活跃KPI
kpis = db.query(KPIDefinition).filter(
KPIDefinition.status == "active",
KPIDefinition.entity_id == entity_id,
).all()
if not kpis:
return {"entity_id": entity_id, "kpis": [], "layers": {}, "overall": None}
# 获取企业信息
ent = db.query(Entity).filter(Entity.id == entity_id).first()
entity_info = {"id": ent.id, "name": ent.name, "short_name": ent.short_name} if ent else {"id": entity_id}
# 单个KPI评分
kpi_scores = []
for k in kpis:
# 取最新实际值
val_query = db.query(KPIValue).filter(
KPIValue.kpi_id == k.id,
KPIValue.actual_value.isnot(None),
)
if period:
val_query = val_query.filter(KPIValue.period == period)
latest_val = val_query.order_by(KPIValue.period.desc()).first()
current_val = latest_val.actual_value if latest_val else None
score, status = _calc_five_tier_score(current_val, k.target_value, is_reverse=(k.kpi_code in REVERSE_INDICATORS))
kpi_scores.append({
"kpi_id": k.id,
"kpi_code": k.kpi_code,
"kpi_name": k.kpi_name,
"dimension": k.dimension,
"target_value": k.target_value,
"current_value": current_val,
"score": score,
"status": status,
"unit": k.unit,
"weight": 10,
"period": latest_val.period if latest_val else None,
})
# BSC四层汇总
layer_map = {
"finance": {"label": "财务", "order": 0},
"customer": {"label": "客户", "order": 1},
"process": {"label": "流程", "order": 2},
"learning": {"label": "学习成长", "order": 3},
}
layers = {}
total_weighted_score = 0
total_weight = 0
for dim_key, dim_info in layer_map.items():
layer_kpis = [s for s in kpi_scores if s["dimension"] == dim_key and s["score"] is not None]
if not layer_kpis:
layers[dim_key] = {"label": dim_info["label"], "score": None, "status": "info", "kpi_count": 0, "weighted_score": None}
continue
w = sum(k["weight"] for k in layer_kpis)
ws = sum(k["score"] * k["weight"] for k in layer_kpis)
avg_score = ws / w if w > 0 else None
avg_status = "success" if avg_score and avg_score >= 4 else ("warning" if avg_score and avg_score >= 3 else "danger") if avg_score else "info"
layers[dim_key] = {
"label": dim_info["label"],
"score": round(avg_score, 2) if avg_score else None,
"status": avg_status,
"kpi_count": len(layer_kpis),
"weighted_score": round(avg_score, 2) if avg_score else None,
}
if avg_score:
total_weighted_score += avg_score * len(layer_kpis)
total_weight += len(layer_kpis)
# 综合得分
overall_score = round(total_weighted_score / total_weight, 2) if total_weight > 0 else None
overall_status = "success" if overall_score and overall_score >= 4 else ("warning" if overall_score and overall_score >= 3 else "danger") if overall_score else "info"
return {
"entity": entity_info,
"kpis": kpi_scores,
"layers": layers,
"overall": {"score": overall_score, "status": overall_status},
}
# ============================================================
# KPI-glossary: 知识资产化 — KPI字典实时加载(供ChatBI财务Bot调用)
# ============================================================
@router.get("/glossary")
def get_kpi_glossary(
entity_id: int = Query(1, ge=1),
db: Session = Depends(get_db),
current_user = Depends(require_auth),
):
"""KPI字典实时加载 — 返回所有KPI的定义、当前值、目标值、公式、维度、阈值
供ChatBI财务Bot在分析前调用,确保口径与系统一致。
返回字段: kpi_code, kpi_name, current_value, target_value, formula, dimension, threshold
"""
kpis = db.query(KPIDefinition).filter(
KPIDefinition.status == "active",
KPIDefinition.entity_id == entity_id,
).order_by(KPIDefinition.kpi_code).all()
result = []
for k in kpis:
# 获取最新实际值
latest_val = db.query(KPIValue).filter(
KPIValue.kpi_id == k.id,
KPIValue.actual_value.isnot(None),
).order_by(KPIValue.period.desc()).first()
current_value = latest_val.actual_value if latest_val else None
latest_period = latest_val.period if latest_val else None
# 组装阈值描述
threshold = None
if k.threshold_green or k.threshold_yellow or k.threshold_red:
parts = []
if k.threshold_green:
parts.append(f"绿灯:{k.threshold_green}")
if k.threshold_yellow:
parts.append(f"黄灯:{k.threshold_yellow}")
if k.threshold_red:
parts.append(f"红灯:{k.threshold_red}")
threshold = " | ".join(parts)
result.append({
"kpi_id": k.id,
"kpi_code": k.kpi_code,
"kpi_name": k.kpi_name,
"dimension": k.dimension,
"category": k.category,
"formula": k.formula,
"formula_desc": k.formula_desc,
"unit": k.unit,
"target_value": k.target_value,
"current_value": current_value,
"latest_period": latest_period,
"threshold": threshold,
"responsible_dept": k.responsible_dept,
"responsible_user": k.responsible_user,
"data_source": k.data_source,
"data_owner": k.data_owner,
"frequency": k.frequency,
"status": k.status,
})
return {
"entity_id": entity_id,
"total": len(result),
"glossary": result,
}
# ============================================================
# KPI-6: KPI三级分解树
# ============================================================
@router.get("/hierarchy")
def get_kpi_hierarchy(
entity_id: int = Query(1, ge=1),
kpi_id: Optional[int] = None,
db: Session = Depends(get_db),
current_user = Depends(require_auth),
):
"""KPI三级分解树:公司→部门→个人"""
query = db.query(KPIHierarchy).filter(KPIHierarchy.entity_id == entity_id)
if kpi_id is not None:
query = query.filter(
(KPIHierarchy.parent_kpi_id == kpi_id) | (KPIHierarchy.child_kpi_id == kpi_id)
)
relations = query.order_by(KPIHierarchy.level).all()
if not relations:
# 无层级数据,返回公司级KPI作为根节点
kpis = db.query(KPIDefinition).filter(
KPIDefinition.status == "active",
KPIDefinition.entity_id == entity_id,
).limit(20).all()
return {
"entity_id": entity_id,
"tree": [{"id": k.id, "kpi_code": k.kpi_code, "kpi_name": k.kpi_name,
"dimension": k.dimension, "level": 1, "children": []} for k in kpis],
"total": len(kpis),
}
# 构建树
kpi_ids = set()
for r in relations:
kpi_ids.add(r.parent_kpi_id)
kpi_ids.add(r.child_kpi_id)
kpi_map = {}
for kid in kpi_ids:
k = db.query(KPIDefinition).filter(KPIDefinition.id == kid).first()
if k:
kpi_map[kid] = {"id": k.id, "kpi_code": k.kpi_code, "kpi_name": k.kpi_name,
"dimension": k.dimension, "level": None, "children": []}
# 分配层级
for r in relations:
if r.parent_kpi_id in kpi_map:
kpi_map[r.parent_kpi_id]["level"] = 1 # 公司级
if r.child_kpi_id in kpi_map:
current_level = kpi_map[r.child_kpi_id].get("level")
new_level = r.level or 2
if current_level is None or current_level > new_level:
kpi_map[r.child_kpi_id]["level"] = new_level
# 构造父子关系
tree = []
added = set()
for r in relations:
parent = kpi_map.get(r.parent_kpi_id)
child = kpi_map.get(r.child_kpi_id)
if parent and child:
child_node = dict(child)
child_node["weight"] = r.weight
child_node["child_name"] = r.child_name
# 避免重复添加
child_key = r.child_kpi_id
existing_child = next(
(c for c in parent["children"] if c["id"] == child_key), None
)
if not existing_child:
parent["children"].append(child_node)
# 收集顶级节点(有子节点且未被引用的parent)
all_child_ids = {r.child_kpi_id for r in relations}
for r in relations:
pid = r.parent_kpi_id
if pid not in all_child_ids or pid == (kpi_id if kpi_id else -1):
if pid not in added and pid in kpi_map:
tree.append(kpi_map[pid])
added.add(pid)
# 如果kpi_id指定,返回该节点为根的子树
if kpi_id is not None and kpi_id in kpi_map:
root = kpi_map[kpi_id]
return {"entity_id": entity_id, "tree": [root], "total": len(tree)}
# 否则按level排序
tree.sort(key=lambda n: (n.get("level") or 99, n["kpi_code"]))
return {"entity_id": entity_id, "tree": tree, "total": len(tree)}
# ============================================================
# KPI-8: KPI因果链追踪
# ============================================================
@router.get("/{kpi_id}/causality-chain")
def get_kpi_causality_chain(
kpi_id: int,
db: Session = Depends(get_db),
current_user = Depends(require_auth),
):
"""KPI因果链追踪 — 返回单个KPI的上下游因果链"""
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
# 上游(驱动当前KPI的因子)
upstream = db.query(KPICausality).filter(KPICausality.target_kpi_id == kpi_id).all()
upstream_list = []
for c in upstream:
src = db.query(KPIDefinition).filter(KPIDefinition.id == c.source_kpi_id).first()
if src:
upstream_list.append({
"causality_id": c.id,
"kpi_id": src.id,
"kpi_code": src.kpi_code,
"kpi_name": src.kpi_name,
"dimension": src.dimension,
"layer": src.dimension,
"strength": c.strength,
"lag_months": c.lag_months,
"direction": c.direction,
"formula": c.formula,
})
# 下游(当前KPI影响的指标)
downstream = db.query(KPICausality).filter(KPICausality.source_kpi_id == kpi_id).all()
downstream_list = []
for c in downstream:
tgt = db.query(KPIDefinition).filter(KPIDefinition.id == c.target_kpi_id).first()
if tgt:
downstream_list.append({
"causality_id": c.id,
"kpi_id": tgt.id,
"kpi_code": tgt.kpi_code,
"kpi_name": tgt.kpi_name,
"dimension": tgt.dimension,
"layer": tgt.dimension,
"strength": c.strength,
"lag_months": c.lag_months,
"direction": c.direction,
"formula": c.formula,
})
return {
"kpi": {
"id": kpi.id,
"kpi_code": kpi.kpi_code,
"kpi_name": kpi.kpi_name,
"dimension": kpi.dimension,
"layer": kpi.dimension,
},
"drives": downstream_list,
"driven_by": upstream_list,
"total_upstream": len(upstream_list),
"total_downstream": len(downstream_list),
}
# ============================================================
# 动态路由(必须在静态路由之后)
# ============================================================
@router.get("/{kpi_id}")
def get_kpi(kpi_id: int, db: Session = Depends(get_db)):
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
return kpi_to_dict(kpi)
def _validate_kpi_data(data: dict, is_update: bool = False):
"""数据治理:入库必检 + 元数据校验"""
errors = []
# 规则1: target_value 必填
tv = data.get("target_value")
if tv is None or (isinstance(tv, (int, float)) and tv < 0 and not is_update):
if not is_update or "target_value" in data:
if tv is None:
errors.append("目标值(target_value)不能为空")
# 规则1: unit 必填
unit = data.get("unit")
if not unit or (isinstance(unit, str) and unit.strip() == ""):
if not is_update or "unit" in data:
errors.append("单位(unit)不能为空")
# 规则2: 元数据必填 — formula/data_source/data_owner
for field, label in [("formula", "计算公式"), ("data_source", "数据来源"), ("data_owner", "数据责任人")]:
val = data.get(field)
if not val or (isinstance(val, str) and val.strip() == ""):
if not is_update or field in data:
errors.append(f"元数据字段'{label}'({field})不能为空")
return errors
@router.post("")
def create_kpi(data: dict, db: Session = Depends(get_db), user=WRITE_ROLES):
# 检查编码唯一性
existing = db.query(KPIDefinition).filter(KPIDefinition.kpi_code == data.get("kpi_code", "")).first()
if existing:
raise HTTPException(400, f"KPI编码 {data['kpi_code']} 已存在")
# 数据治理校验
errs = _validate_kpi_data(data, is_update=False)
if errs:
raise HTTPException(422, detail={"message": "数据校验不通过", "errors": errs})
kpi = KPIDefinition(**data)
db.add(kpi)
db.commit()
db.refresh(kpi)
_log(db, 1, "create", "kpi", kpi.id, data)
return kpi_to_dict(kpi)
@router.put("/{kpi_id}")
def update_kpi(kpi_id: int, data: dict, db: Session = Depends(get_db), user=WRITE_ROLES):
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
# 数据治理校验(更新时只检查传了但为空的字段)
errs = _validate_kpi_data(data, is_update=True)
if errs:
raise HTTPException(422, detail={"message": "数据校验不通过", "errors": errs})
for k, v in data.items():
if hasattr(kpi, k) and v is not None:
setattr(kpi, k, v)
db.commit()
_log(db, 1, "update", "kpi", kpi_id, data)
return kpi_to_dict(kpi)
@router.delete("/{kpi_id}")
def delete_kpi(kpi_id: int, db: Session = Depends(get_db), user=WRITE_ROLES):
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if kpi:
kpi.status = "disabled"
db.commit()
return {"message": "已删除"}
@router.put("/{kpi_id}/restore")
def restore_kpi(kpi_id: int, db: Session = Depends(get_db), user=WRITE_ROLES):
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if kpi:
kpi.status = "active"
db.commit()
return {"message": "已恢复"}
def kpi_to_dict(k):
d = {c.name: getattr(k, c.name) for c in k.__table__.columns}
# 附加战略地图信息
if k.map_id:
from app.database import get_session_local
try:
sess = get_session_local()()
m = sess.query(StrategicMap).filter(StrategicMap.id == k.map_id).first()
d["map_title"] = m.title if m else None
sess.close()
except:
d["map_title"] = None
else:
d["map_title"] = None
return d
@router.put("/{kpi_id}/associate-map")
def associate_kpi_map(kpi_id: int, data: dict, db: Session = Depends(get_db), user=WRITE_ROLES):
"""关联KPI到战略地图"""
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
map_id = data.get("map_id")
if map_id is not None:
m = db.query(StrategicMap).filter(StrategicMap.id == map_id).first()
if not m:
raise HTTPException(404, "战略地图不存在")
kpi.map_id = map_id
db.commit()
_log(db, 1, "update", "kpi", kpi_id, {"action": "associate-map", "map_id": map_id})
return kpi_to_dict(kpi)
def _log(db, user_id, action, target_type, target_id, detail):
log = OperationLog(user_id=user_id, action=action, target_type=target_type, target_id=target_id, detail=json.dumps(detail, ensure_ascii=False) if detail else None)
db.add(log)
db.commit()
@router.get("/{kpi_id}/objectives")
def get_kpi_objectives(kpi_id: int, db: Session = Depends(get_db)):
"""查看KPI所属的目标和战略地图"""
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
# 通过 kpi_definitions.objective 字段关联目标
# 也通过 map_id 关联地图
result = {
"kpi": {"id": kpi.id, "kpi_code": kpi.kpi_code, "kpi_name": kpi.kpi_name},
"objectives": [],
"map": None,
}
if kpi.map_id:
m = db.query(StrategicMap).filter(StrategicMap.id == kpi.map_id).first()
if m:
result["map"] = {"id": m.id, "title": m.title, "status": m.status}
if kpi.objective:
objs = db.query(MapObjective).filter(
MapObjective.map_id == kpi.map_id,
MapObjective.name == kpi.objective,
).all()
result["objectives"] = [{"id": o.id, "name": o.name, "dimension_key": o.dimension_key} for o in objs]
return result