Commit 9de8abeb authored by Data Governance Dev's avatar Data Governance Dev

fix(step4): 修复 KeyError(column_name/column_comment) 和 TypeError(lambda kwargs);移除 LLM 依赖改为纯程序化检查

parent a9fe801f
......@@ -69,21 +69,29 @@ async def llm_status():
async def list_steps():
return StepsListResponse(steps=[
StepInfo(num=1, title="获取数据字典",
description="从 information_schema 抓取所有表/字段元数据", requires_db=True),
description="从 information_schema 抓取所有表/字段元数据",
requires_db=True, llm_mode="none"),
StepInfo(num=2, title="表合并与冗余字段分析",
description="离线分析:找结构高度相似的表 + 高频字段(LLM 优化建议)", requires_db=False, llm_enhanced=True),
description="离线:找结构高度相似的表 + 高频字段(LLM 必须:合并 verdict 与冗余分类)",
requires_db=False, llm_mode="required"),
StepInfo(num=3, title="数据验证",
description="连库验证:合并候选实际行数 + 行政区划孤儿 + 数据质量问题", requires_db=True),
description="连库验证:合并候选实际行数 + 行政区划孤儿 + 数据质量问题",
requires_db=True, llm_mode="none"),
StepInfo(num=4, title="大范围空字段扫描",
description="单表全列扫,统计 NULL/空值比例(LLM 分类成因)", requires_db=True, llm_enhanced=True),
description="单表全列扫,统计 NULL/空值比例(纯程序化检查,不依赖 LLM)",
requires_db=True, llm_mode="none"),
StepInfo(num=5, title="缺失注释字段 + 推测",
description="离线:找无注释字段,规则+LLM 推测语义", requires_db=False, llm_enhanced=True),
description="离线:找无注释字段,LLM 必须:批量推测语义",
requires_db=False, llm_mode="required"),
StepInfo(num=6, title="字段长度检查",
description="离线:识别字段定义长度超出国家标准(LLM 推断自定义字段)", requires_db=False, llm_enhanced=True),
description="离线:识别字段定义长度超出国家标准(LLM 可选:自定义字段推断)",
requires_db=False, llm_mode="optional"),
StepInfo(num=7, title="国家标准校验",
description="连库抽样:身份证/USCC/手机号/行政区划符合性", requires_db=True),
description="连库抽样:身份证/USCC/手机号/行政区划符合性",
requires_db=True, llm_mode="none"),
StepInfo(num=8, title="报告生成",
description="汇总所有 findings 生成 Markdown + Word 报告", requires_db=False),
description="汇总所有 findings 生成 Markdown + Word 报告",
requires_db=False, llm_mode="none"),
])
......@@ -141,7 +149,8 @@ async def connect_test(req: TestConnectionRequest):
async def create_job(req: ConnectRequest):
logger.info(
f"POST /api/jobs db={req.database}, db_type={req.db_type}, "
f"steps={req.steps or 'all'}, enable_llm={req.enable_llm}"
f"steps={req.steps or 'all'}, enable_llm={req.enable_llm} "
f"(注:enable_llm=False 时,Step 2/5 会因 LLM-required 而失败)"
)
jm = JobManager.get_instance()
job = jm.submit(req)
......@@ -204,18 +213,36 @@ async def get_job_result(job_id: str):
# ── 日志流(SSE) ──
@router.get("/jobs/{job_id}/logs/stream")
async def stream_logs(job_id: str):
async def stream_logs(job_id: str, request: Request):
from sse_starlette.sse import EventSourceResponse
jm = JobManager.get_instance()
job = jm.get(job_id)
if not job:
logger.warning(f"GET /api/jobs/{job_id}/logs/stream 任务不存在")
raise HTTPException(status_code=404, detail="任务不存在")
logger.info(f"GET /api/jobs/{job_id}/logs/stream SSE 订阅开始 (历史 {len(job.logs.snapshot())} 条)")
# 浏览器 EventSource 在自动重连时会发送 Last-Event-ID header,
# 后端据此跳过已发条目,避免日志被重复 append 到前端列表
last_event_id = request.headers.get("Last-Event-ID", "")
last_seq = 0
if last_event_id:
try:
last_seq = int(last_event_id)
except ValueError:
last_seq = 0
logger.info(
f"GET /api/jobs/{job_id}/logs/stream SSE 订阅开始 "
f"(历史 {len(job.logs.snapshot())} 条, last_seq={last_seq})"
)
async def event_generator():
async for entry in jm.stream_logs(job_id):
yield {"event": "log", "data": _json_dumps(entry)}
async for entry in jm.stream_logs(job_id, last_seq=last_seq):
# 设置 SSE event id —— 浏览器存下来,重连时通过 Last-Event-ID 回传
eid = entry.get("_seq")
payload = {"event": "log", "data": _json_dumps(entry)}
if eid is not None:
payload["id"] = str(eid)
yield payload
return EventSourceResponse(event_generator())
......
......@@ -32,18 +32,29 @@ logger = logging.getLogger(__name__)
# ── Step 注册表 ────────────────────────────────────────────
# 顺序敏感:前面 Step 的产出可被后面 Step 复用
STEP_REGISTRY: list[tuple[int, str, Callable, bool, bool]] = [
# (num, title, fn, requires_db, llm_enhanced)
(1, "获取数据字典", "_run_step1", True, False),
(2, "表合并与冗余字段分析", "_run_step2", False, True),
(3, "数据验证", "_run_step3", True, False),
(4, "大范围空字段扫描", "_run_step4", True, True),
(5, "缺失注释检查 + 推测", "_run_step5", False, True),
(6, "字段长度检查", "_run_step6", False, True),
(7, "国家标准校验", "_run_step7", True, False),
(8, "报告生成", "_run_step8", False, False),
#
# llm_mode 取值:
# "none" —— 不用 LLM
# "optional" —— LLM 可用就用,失败降级到规则推理(Step 4、6)
# "required" —— 必须 LLM;缺 LLM = 任务失败(Step 2、5)
STEP_REGISTRY: list[tuple[int, str, Callable, bool, str]] = [
# (num, title, fn, requires_db, llm_mode)
(1, "获取数据字典", "_run_step1", True, "none"),
(2, "表合并与冗余字段分析", "_run_step2", False, "required"),
(3, "数据验证", "_run_step3", True, "none"),
(4, "大范围空字段扫描", "_run_step4", True, "none"),
(5, "缺失注释检查 + 推测", "_run_step5", False, "required"),
(6, "字段长度检查", "_run_step6", False, "optional"),
(7, "国家标准校验", "_run_step7", True, "none"),
(8, "报告生成", "_run_step8", False, "none"),
]
# ── 关键 Step 集合 ────────────────────────────────────────
# 失败会导致任务标 JOB_FAILED,不返回部分结果
# Step 1: 后续所有 Step 依赖数据字典
# Step 2 / 5: LLM-required,缺 LLM 视为关键失败
CRITICAL_STEPS = {1, 2, 5}
# ── 工具函数 ────────────────────────────────────────────
def _check_cancel(event: asyncio.Event) -> None:
......@@ -95,7 +106,7 @@ async def run_governance_workflow(
_prev_web_level = web_logger.level
# 关键 Step 失败应直接终止流程(其他 Step 依赖其产出)
CRITICAL_STEPS = {1}
# CRITICAL_STEPS 在模块顶部定义;此处不再重复
llm = get_llm_client() if enable_llm else None
......@@ -111,9 +122,30 @@ async def run_governance_workflow(
# ── 执行每个 Step ──
log("INFO", f"本次治理配置: db_type={cfg.db_type}, host={cfg.host}:{cfg.port}, db={cfg.database}")
log("INFO", f"计划执行步骤: {steps}")
total_steps = len([n for n, *_ in STEP_REGISTRY if n in steps])
for step_idx, (num, title, fn_name, requires_db, llm_enhanced) in enumerate(STEP_REGISTRY, 1):
total_steps = len([n for n, _, _, _, _ in STEP_REGISTRY if n in steps])
# ── 启动预检:LLM-required 步骤 + LLM 不可用 → fail-fast ──
required_planned = [
n for n, _, _, _, mode in STEP_REGISTRY
if mode == "required" and n in steps
]
llm_available = bool(llm and llm.available)
if required_planned and not llm_available:
msg = (
f"用户计划执行 LLM-required 步骤 {required_planned},"
f"但 LLM 不可用(未配置 API Key 或客户端初始化失败)。"
f"请在 web/configs/llm.yaml 配置 api_key,"
f"或在提交任务时去掉这些步骤后重试。"
)
log("ERROR", msg)
logger.error(msg)
raise RuntimeError(msg)
if required_planned and llm_available:
log("INFO",
f"LLM-required 步骤 {required_planned} 已就绪 "
f"(provider={llm.cfg.provider}, model={llm.cfg.effective_model})")
for step_idx, (num, title, fn_name, requires_db, llm_mode) in enumerate(STEP_REGISTRY, 1):
if num not in steps:
continue
_check_cancel(cancel_event)
......@@ -125,15 +157,18 @@ async def run_governance_workflow(
log("INFO", " 需要连接数据库", step=str(num))
step_started = time.monotonic()
# 标记 LLM 状态
llm_note = ""
if llm_enhanced:
# 标记 LLM 状态(按 mode 分级提示)
if llm_mode == "required":
log("INFO",
f" ⚠ LLM 必须: provider={llm.cfg.provider}, model={llm.cfg.effective_model}",
step=str(num))
elif llm_mode == "optional":
if llm and llm.available:
llm_note = f"(LLM 增强已启用: provider={llm.cfg.provider}, model={llm.cfg.effective_model})"
log("INFO",
f" · LLM 增强已启用: provider={llm.cfg.provider}, model={llm.cfg.effective_model}",
step=str(num))
else:
llm_note = "(LLM 未配置,将使用规则推理)"
if llm_note:
log("INFO", llm_note, step=str(num))
log("INFO", " · LLM 未配置,将使用规则推理", step=str(num))
try:
# 在线程池里跑同步步骤(DB 查询 / 文件 IO)
......@@ -143,7 +178,7 @@ async def run_governance_workflow(
lambda: step_fn(
cfg=cfg,
step_outputs=step_outputs,
llm=llm if llm_enhanced else None,
llm=llm if llm_mode != "none" else None,
findings_dir=findings_dir,
log=log,
cancel_event=cancel_event,
......@@ -270,9 +305,9 @@ def _run_step3(cfg, step_outputs, llm, findings_dir, log, cancel_event):
def _run_step4(cfg, step_outputs, llm, findings_dir, log, cancel_event):
"""Step 4: 空字段扫描(连库)"""
"""Step 4: 空字段扫描(连库)——纯程序化检查,不依赖 LLM"""
from .step_impl.step4_empty_fields import run_step4
data = run_step4(cfg, log=log, llm=llm)
data = run_step4(cfg, log=log)
return {"section_key": "empty_fields", "data": data}
......
"""Step 4: 大范围空字段扫描(连库 + LLM 增强)
"""Step 4: 大范围空字段扫描(纯程序化检查)
对每张表跑动态 SQL 统计:
- COUNT(*) 总行数
......@@ -8,7 +8,7 @@
- 空值率 ≥80% → "high" 高空
- 空值率 ≥50% → "mid" 中空
LLM 增强(可选):对 top 20 高空字段分类「业务不需要 / 采集漏了 / 预留字段」。
不依赖 LLM,完全基于数据字典做空字段占比检查。
SQL 全部走 web/sql/empty_fields/ 模板,Python 端只构造列子句(empty_clauses)。
"""
......@@ -20,26 +20,17 @@ from collections import defaultdict
from typing import Callable
from ..db_adapter import DBConfig, open_db, quote_ident, TEXT_TYPES
from ..llm import LLMClient
from web.sql.loader import get_sql_loader
logger = logging.getLogger(__name__)
# 字段名 → 业务关键性提示(帮助 LLM 更准确分类)
CRITICAL_FIELDS = {
"id", "name", "code", "type", "status", "user_id", "create_time", "update_time",
"amount", "price", "total", "count", "project_id", "user_name",
}
def run_step4(cfg: DBConfig, log: Callable | None = None,
llm: LLMClient | None = None,
min_rows_to_check: int = 5,
high_threshold: float = 0.80,
mid_threshold: float = 0.50,
sample_size: int = 3) -> dict:
"""扫描所有表的空字段情况"""
"""扫描所有表的空字段情况(纯程序化,不依赖 LLM)"""
if log:
log("INFO",
f"Step 4 开始扫描空字段 (high≥{high_threshold:.0%}, mid≥{mid_threshold:.0%}, min_rows={min_rows_to_check})",
......@@ -47,7 +38,7 @@ def run_step4(cfg: DBConfig, log: Callable | None = None,
# Step 1 产出已在 step_outputs 里,通过 cfg 重新拉一遍元数据
from .step1_data_dict import run_step1
dict_data = run_step1(cfg, log=lambda lvl, msg: log(lvl, msg, step="4") if log else None)
dict_data = run_step1(cfg, log=lambda lvl, msg, **_kw: log(lvl, msg, step="4") if log else None)
columns = dict_data["data_dictionary"]
table_summary = dict_data["table_summary"]
......@@ -128,10 +119,10 @@ def run_step4(cfg: DBConfig, log: Callable | None = None,
field_record = {
"table_name": table,
"table_comment": col.get("table_comment", ""),
"column_name": col["name"],
"column_name": col["column_name"],
"data_type": col["data_type"],
"column_type": col["column_type"],
"column_comment": col.get("comment", ""),
"column_comment": col.get("column_comment", ""),
"total_rows": row_count,
"empty_count": empty_count,
"empty_rate": round(empty_rate, 4),
......@@ -189,41 +180,6 @@ def run_step4(cfg: DBConfig, log: Callable | None = None,
f"跳过 {len(skipped)} 张",
step="4")
# LLM 增强:对 top 20 高空字段分类
if llm and llm.available and high_fields_all:
if log:
log("INFO", f"调用 LLM 分类 top {min(20, len(high_fields_all))} 个高空字段成因", step="4")
classified = []
for idx, fld in enumerate(high_fields_all[:20], 1):
try:
# 抽样非空样本(这里简化为空,仅展示结构)
r = llm.classify_empty_field(
table_name=fld["table_name"],
table_comment=fld["table_comment"],
column_name=fld["column_name"],
column_comment=fld["column_comment"],
data_type=fld["data_type"],
empty_rate=fld["empty_rate"],
sample_values=[],
)
if r:
fld["llm_classification"] = r.get("classification", "unknown")
fld["llm_reasoning"] = r.get("reasoning", "")
fld["llm_action"] = r.get("action", "")
classified.append(fld)
if log:
log("DEBUG",
f" · [{idx}/{min(20, len(high_fields_all))}] "
f"{fld['table_name']}.{fld['column_name']} → {r.get('classification')}",
step="4")
except Exception as e:
if log:
log("WARN", f"LLM 分类失败 ({fld['table_name']}.{fld['column_name']}): {e}", step="4")
summary["llm_classified_count"] = len(classified)
summary["llm_classified_sample"] = classified[:10]
if log:
log("INFO", f"LLM 分类成功 {len(classified)}/{min(20, len(high_fields_all))}", step="4")
return {
"summary": summary,
"per_table": per_table,
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment