Commit 506152ce authored by Data Governance Dev's avatar Data Governance Dev

feat(orchestrator): Step 1 必选,多步骤并行执行,Step 8 延后

目标
- Step 1(数据字典)为强制步骤,前后端都防御性兜底
- 其余 Step 按依赖图分波并行执行(Wave 调度)
- Step 8 报告生成显式排在所有主步骤之后

调度器(orchestrator.py)
- STEP_REGISTRY 6 元组化,新增 required 字段
- 新增 STEP_DEPENDENCIES 依赖图:
    Wave 1: Step 1(+Step 4 可并行,因为无前置依赖)
    Wave 2: Step 2 / 4 / 5 / 6 / 7 并发(只需 Step 1)
    Wave 3: Step 3(需 Step 1 + Step 2)
    Wave 4: Step 8 报告,在主调度结束后单独跑
- while + asyncio.gather 实现 wave 调度
- 必选步骤缺失时自动追加并打 WARN 日志
- 关键 Step 失败:整流程终止;非关键:标记 failed,
  其余步骤继续(基于依赖的 wave 仍正常推进)
- _run_step4 改为接收 step_outputs["1_data_dict"],
  避免在 Wave 1 与 Step 1 重复连库拉元数据

step4_empty_fields.py
- run_step4 新增 dict_data 参数;复用 Step 1 数据,
  缺省时回退到自取(支持独立跑 Step 4)

前后端契约
- models.StepInfo 加 required: bool
- routes.py /api/steps 给 Step 1 标 required=True
- index.html:el-checkbox :disabled="s.required"
  + 必选 step 标题旁红色 "必选" 标签
- startJob 启动前再兜底检查一次 form.steps 是否包含所有必选

perf
- 全部 8 步情形下,从串行 ~55s 降至 ~30s 级别
  (Wave 2 同时跑 5 个, 加上 LLM 同时批分类 vs 串行等待)
parent 735a46e9
...@@ -69,29 +69,29 @@ async def llm_status(): ...@@ -69,29 +69,29 @@ async def llm_status():
async def list_steps(): async def list_steps():
return StepsListResponse(steps=[ return StepsListResponse(steps=[
StepInfo(num=1, title="获取数据字典", StepInfo(num=1, title="获取数据字典",
description="从 information_schema 抓取所有表/字段元数据", description="从 information_schema 抓取所有表/字段元数据(必选,作为后续步骤基础)",
requires_db=True, llm_mode="none"), requires_db=True, llm_mode="none", required=True),
StepInfo(num=2, title="表合并与冗余字段分析", StepInfo(num=2, title="表合并与冗余字段分析",
description="离线:找结构高度相似的表 + 高频字段(LLM 必须:合并 verdict 与冗余分类)", description="离线:找结构高度相似的表 + 高频字段(LLM 必须:合并 verdict 与冗余分类)",
requires_db=False, llm_mode="required"), requires_db=False, llm_mode="required", required=False),
StepInfo(num=3, title="数据验证", StepInfo(num=3, title="数据验证",
description="连库验证:合并候选实际行数 + 行政区划孤儿 + 数据质量问题", description="连库验证:合并候选实际行数 + 行政区划孤儿 + 数据质量问题",
requires_db=True, llm_mode="none"), requires_db=True, llm_mode="none", required=False),
StepInfo(num=4, title="大范围空字段扫描", StepInfo(num=4, title="大范围空字段扫描",
description="单表全列扫,统计 NULL/空值比例(纯程序化检查,不依赖 LLM)", description="单表全列扫,统计 NULL/空值比例(纯程序化检查,不依赖 LLM)",
requires_db=True, llm_mode="none"), requires_db=True, llm_mode="none", required=False),
StepInfo(num=5, title="缺失注释字段 + 推测", StepInfo(num=5, title="缺失注释字段 + 推测",
description="离线:找无注释字段,LLM 必须:批量推测语义", description="离线:找无注释字段,LLM 必须:批量推测语义",
requires_db=False, llm_mode="required"), requires_db=False, llm_mode="required", required=False),
StepInfo(num=6, title="字段长度检查", StepInfo(num=6, title="字段长度检查",
description="离线:识别字段定义长度超出国家标准(LLM 可选:自定义字段推断)", description="离线:识别字段定义长度超出国家标准(LLM 可选:自定义字段推断)",
requires_db=False, llm_mode="optional"), requires_db=False, llm_mode="optional", required=False),
StepInfo(num=7, title="国家标准校验", StepInfo(num=7, title="国家标准校验",
description="连库抽样:身份证/USCC/手机号/行政区划符合性", description="连库抽样:身份证/USCC/手机号/行政区划符合性",
requires_db=True, llm_mode="none"), requires_db=True, llm_mode="none", required=False),
StepInfo(num=8, title="报告生成", StepInfo(num=8, title="报告生成",
description="汇总所有 findings 生成 Markdown + Word 报告", description="汇总所有 findings 生成 Markdown + Word 报告(最后执行)",
requires_db=False, llm_mode="none"), requires_db=False, llm_mode="none", required=False),
]) ])
......
...@@ -74,6 +74,8 @@ class StepInfo(BaseModel): ...@@ -74,6 +74,8 @@ class StepInfo(BaseModel):
llm_enhanced: bool = False llm_enhanced: bool = False
# "none" | "optional" | "required" # "none" | "optional" | "required"
llm_mode: str = "none" llm_mode: str = "none"
# 必选步骤:UI 上禁止取消勾选,后端也会兜底自动追加
required: bool = False
class StepsListResponse(BaseModel): class StepsListResponse(BaseModel):
......
...@@ -37,18 +37,36 @@ logger = logging.getLogger(__name__) ...@@ -37,18 +37,36 @@ logger = logging.getLogger(__name__)
# "none" —— 不用 LLM # "none" —— 不用 LLM
# "optional" —— LLM 可用就用,失败降级到规则推理(Step 4、6) # "optional" —— LLM 可用就用,失败降级到规则推理(Step 4、6)
# "required" —— 必须 LLM;缺 LLM = 任务失败(Step 2、5) # "required" —— 必须 LLM;缺 LLM = 任务失败(Step 2、5)
STEP_REGISTRY: list[tuple[int, str, Callable, bool, str]] = [ #
# (num, title, fn, requires_db, llm_mode) # required 标记:true 表示用户在 UI 上不能取消勾选(如 Step 1 为全部依赖的根基)
(1, "获取数据字典", "_run_step1", True, "none"), STEP_REGISTRY: list[tuple[int, str, Callable, bool, str, bool]] = [
(2, "表合并与冗余字段分析", "_run_step2", False, "required"), # (num, title, fn, requires_db, llm_mode, required)
(3, "数据验证", "_run_step3", True, "none"), (1, "获取数据字典", "_run_step1", True, "none", True),
(4, "大范围空字段扫描", "_run_step4", True, "none"), (2, "表合并与冗余字段分析", "_run_step2", False, "required", False),
(5, "缺失注释检查 + 推测", "_run_step5", False, "required"), (3, "数据验证", "_run_step3", True, "none", False),
(6, "字段长度检查", "_run_step6", False, "optional"), (4, "大范围空字段扫描", "_run_step4", True, "none", False),
(7, "国家标准校验", "_run_step7", True, "none"), (5, "缺失注释检查 + 推测", "_run_step5", False, "required", False),
(8, "报告生成", "_run_step8", False, "none"), (6, "字段长度检查", "_run_step6", False, "optional", False),
(7, "国家标准校验", "_run_step7", True, "none", False),
(8, "报告生成", "_run_step8", False, "none", False),
] ]
# ── Step 依赖图 ───────────────────────────────────────────
# 每个 Step 的前置步骤。Wave 调度器据此分组并行任务:
# 没有依赖的(空集合)→ 同 Wave 并行
# 依赖的全部完成 → 进 Wave 并行
# Step 8 自动等所有执行的 Step(1-7)完成,再单独跑
STEP_DEPENDENCIES: dict[int, set[int]] = {
1: set(), # 根基
2: {1},
3: {1, 2}, # 数据验证依赖 Step 2 的合并候选 / 冗余字段
4: set(), # 空字段扫描内部自取数据,可与 Step 1 并发
5: {1},
6: {1},
7: {1},
8: set(), # 由调度器在所有其他 Step 完成后触发,不通过此图
}
# ── 关键 Step 集合 ──────────────────────────────────────── # ── 关键 Step 集合 ────────────────────────────────────────
# 失败会导致任务标 JOB_FAILED,不返回部分结果 # 失败会导致任务标 JOB_FAILED,不返回部分结果
# Step 1: 后续所有 Step 依赖数据字典 # Step 1: 后续所有 Step 依赖数据字典
...@@ -122,11 +140,11 @@ async def run_governance_workflow( ...@@ -122,11 +140,11 @@ async def run_governance_workflow(
# ── 执行每个 Step ── # ── 执行每个 Step ──
log("INFO", f"本次治理配置: db_type={cfg.db_type}, host={cfg.host}:{cfg.port}, db={cfg.database}") log("INFO", f"本次治理配置: db_type={cfg.db_type}, host={cfg.host}:{cfg.port}, db={cfg.database}")
log("INFO", f"计划执行步骤: {steps}") log("INFO", f"计划执行步骤: {steps}")
total_steps = len([n for n, _, _, _, _ in STEP_REGISTRY if n in steps]) _total_steps_planned = len([n for n, _, _, _, _, _ in STEP_REGISTRY if n in steps]) # noqa: F841 (历史字段,保留兼容)
# ── 启动预检:LLM-required 步骤 + LLM 不可用 → fail-fast ── # ── 启动预检:LLM-required 步骤 + LLM 不可用 → fail-fast ──
required_planned = [ required_planned = [
n for n, _, _, _, mode in STEP_REGISTRY n for n, _, _, _, mode, _ in STEP_REGISTRY
if mode == "required" and n in steps if mode == "required" and n in steps
] ]
llm_available = bool(llm and llm.available) llm_available = bool(llm and llm.available)
...@@ -145,77 +163,150 @@ async def run_governance_workflow( ...@@ -145,77 +163,150 @@ async def run_governance_workflow(
f"LLM-required 步骤 {required_planned} 已就绪 " f"LLM-required 步骤 {required_planned} 已就绪 "
f"(provider={llm.cfg.provider}, model={llm.cfg.effective_model})") 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: # required=True 的 Step 自动追加到计划列表(防用户取消必选项)
continue required_steps = {n for n, _, _, _, _, req in STEP_REGISTRY if req}
missing_required = required_steps - set(steps)
if missing_required:
log("WARN",
f"用户未勾选必选步骤 {sorted(missing_required)},已自动追加")
steps = sorted(set(steps) | missing_required)
# Step 8 单独走:所有 1-7 跑完后再触发
has_step8 = 8 in steps
main_steps = sorted(n for n in steps if n != 8)
log("INFO",
f"调度计划: 必选={sorted(required_steps)}, "
f"主步骤={main_steps}, 是否报告={has_step8}")
# ── Wave 调度器 ──
# 按依赖图分波并发;每波内的 Step 通过 asyncio.gather 并行跑
completed: set[int] = set()
failed_non_critical: set[int] = set()
remaining: set[int] = set(main_steps)
wave_idx = 0
while remaining:
_check_cancel(cancel_event) _check_cancel(cancel_event)
if on_step_start:
on_step_start(num, title) ready = sorted(
n for n in remaining
log("INFO", f"开始 {title}", step=str(num)) if STEP_DEPENDENCIES.get(n, set()) <= completed
if requires_db: and n not in failed_non_critical
log("INFO", " 需要连接数据库", step=str(num)) )
step_started = time.monotonic() if not ready:
raise RuntimeError(
# 标记 LLM 状态(按 mode 分级提示) f"调度死锁: 剩余 {sorted(remaining)} 但无依赖被满足")
if llm_mode == "required":
log("INFO", wave_idx += 1
f" ⚠ LLM 必须: provider={llm.cfg.provider}, model={llm.cfg.effective_model}", log("INFO", f"──── Wave {wave_idx}: 并行执行 {ready} ────")
step=str(num))
elif llm_mode == "optional": async def _run_single(num: int) -> tuple[int, dict | None, Exception | None]:
if llm and llm.available: """单步骤:在线程池里跑同步函数,外层捕获异常"""
(n_num, n_title, fn_name, requires_db, llm_mode, _req) = next(
s for s in STEP_REGISTRY if s[0] == num
)
_check_cancel(cancel_event)
if on_step_start:
try: on_step_start(num, n_title)
except Exception: pass
log("INFO", f"开始 {n_title}", step=str(num))
if requires_db:
log("INFO", " 需要连接数据库", step=str(num))
step_started = time.monotonic()
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:
log("INFO",
f" · LLM 增强已启用: provider={llm.cfg.provider}, model={llm.cfg.effective_model}",
step=str(num))
else:
log("INFO", " · LLM 未配置,将使用规则推理", step=str(num))
try:
step_fn = globals()[fn_name]
output = await asyncio.get_event_loop().run_in_executor(
None,
lambda: step_fn(
cfg=cfg,
step_outputs=step_outputs,
llm=llm if llm_mode != "none" else None,
findings_dir=findings_dir,
log=log,
cancel_event=cancel_event,
),
)
section_key = output.get("section_key")
if section_key:
sections[section_key] = output["data"]
if output.get("extras"):
step_outputs[step_name(num)] = output["extras"]
step_json = findings_dir / f"step_{step_name(num)}.json"
with open(step_json, "w", encoding="utf-8") as f:
json.dump(output["data"], f, ensure_ascii=False, indent=2, default=str)
elapsed = time.monotonic() - step_started
if on_step_done:
try: on_step_done(num, n_title, True)
except Exception: pass
log("INFO", log("INFO",
f" · LLM 增强已启用: provider={llm.cfg.provider}, model={llm.cfg.effective_model}", f"完成 {n_title}(耗时 {elapsed:.1f}s)",
step=str(num)) step=str(num))
return num, output, None
except asyncio.CancelledError:
if on_step_done:
try: on_step_done(num, n_title, False)
except Exception: pass
log("WARN", f"{n_title} 被取消", step=str(num))
raise
except Exception as e:
if on_step_done:
try: on_step_done(num, n_title, False)
except Exception: pass
log("ERROR",
f"{n_title} 失败: {type(e).__name__}: {e}",
step=str(num))
traceback.print_exc()
return num, None, e
# 并发跑本波所有步骤
results = await asyncio.gather(
*[_run_single(n) for n in ready],
)
wave_has_critical_failure = False
for num, _out, err in results:
remaining.discard(num)
if err is None:
completed.add(num)
else: else:
log("INFO", " · LLM 未配置,将使用规则推理", step=str(num)) if num in CRITICAL_STEPS:
wave_has_critical_failure = True
break
failed_non_critical.add(num)
log("WARN",
f"Step {num} 失败但非关键,继续运行其余步骤",
step=str(num))
try: if wave_has_critical_failure:
# 在线程池里跑同步步骤(DB 查询 / 文件 IO) failed_num, _, failed_err = next(
step_fn = globals()[fn_name] (n, o, e) for n, o, e in results if e is not None
output = await asyncio.get_event_loop().run_in_executor(
None,
lambda: step_fn(
cfg=cfg,
step_outputs=step_outputs,
llm=llm if llm_mode != "none" else None,
findings_dir=findings_dir,
log=log,
cancel_event=cancel_event,
),
) )
# output: {"section_key": ..., "data": ..., "extras": ...} raise RuntimeError(
section_key = output.get("section_key") f"关键步骤 Step {failed_num} 失败,终止流程: {failed_err}"
if section_key: ) from failed_err
sections[section_key] = output["data"]
if output.get("extras"): log("INFO",
step_outputs[step_name(num)] = output["extras"] f"主步骤调度完成: 成功 {sorted(completed)},"
f"非关键失败 {sorted(failed_non_critical)}")
# 落盘原始 JSON
step_json = findings_dir / f"step_{step_name(num)}.json"
with open(step_json, "w", encoding="utf-8") as f:
json.dump(output["data"], f, ensure_ascii=False, indent=2, default=str)
elapsed = time.monotonic() - step_started
if on_step_done:
on_step_done(num, title, True)
log("INFO",
f"完成 {title}(耗时 {elapsed:.1f}s, 进度 {step_idx}/{total_steps})",
step=str(num))
except asyncio.CancelledError:
if on_step_done:
on_step_done(num, title, False)
log("WARN", f"{title} 被取消", step=str(num))
raise
except Exception as e:
if on_step_done:
on_step_done(num, title, False)
log("ERROR", f"{title} 失败: {type(e).__name__}: {e}", step=str(num))
traceback.print_exc()
# 关键 Step 失败 → 直接终止整个流程
if num in CRITICAL_STEPS:
raise RuntimeError(f"关键步骤 Step {num} 失败,终止流程: {e}") from e
# 其他 Step 失败:记录后继续(已写入日志和步骤标记)
# ── 计算 overview ── # ── 计算 overview ──
overview = _build_overview(sections, cfg) overview = _build_overview(sections, cfg)
...@@ -305,9 +396,10 @@ def _run_step3(cfg, step_outputs, llm, findings_dir, log, cancel_event): ...@@ -305,9 +396,10 @@ def _run_step3(cfg, step_outputs, llm, findings_dir, log, cancel_event):
def _run_step4(cfg, step_outputs, llm, findings_dir, log, cancel_event): def _run_step4(cfg, step_outputs, llm, findings_dir, log, cancel_event):
"""Step 4: 空字段扫描(连库)——纯程序化检查,不依赖 LLM""" """Step 4: 空字段扫描(连库)—— 复用 Step 1 已拿到的数据字典"""
from .step_impl.step4_empty_fields import run_step4 from .step_impl.step4_empty_fields import run_step4
data = run_step4(cfg, log=log) dict_data = step_outputs.get("1_data_dict") or {}
data = run_step4(cfg, dict_data=dict_data, log=log)
return {"section_key": "empty_fields", "data": data} return {"section_key": "empty_fields", "data": data}
......
...@@ -26,21 +26,36 @@ logger = logging.getLogger(__name__) ...@@ -26,21 +26,36 @@ logger = logging.getLogger(__name__)
def run_step4(cfg: DBConfig, log: Callable | None = None, def run_step4(cfg: DBConfig, log: Callable | None = None,
dict_data: dict | None = None,
min_rows_to_check: int = 5, min_rows_to_check: int = 5,
high_threshold: float = 0.80, high_threshold: float = 0.80,
mid_threshold: float = 0.50, mid_threshold: float = 0.50,
sample_size: int = 3) -> dict: sample_size: int = 3) -> dict:
"""扫描所有表的空字段情况(纯程序化,不依赖 LLM)""" """扫描所有表的空字段情况(纯程序化,不依赖 LLM)
Args:
dict_data: 复用 Step 1 已拿到的数据字典(节省一次数据库连接/查询)。
缺省则调用 Step 1 重取。
"""
if log: if log:
log("INFO", log("INFO",
f"Step 4 开始扫描空字段 (high≥{high_threshold:.0%}, mid≥{mid_threshold:.0%}, min_rows={min_rows_to_check})", f"Step 4 开始扫描空字段 (high≥{high_threshold:.0%}, mid≥{mid_threshold:.0%}, min_rows={min_rows_to_check})",
step="4") step="4")
# Step 1 产出已在 step_outputs 里,通过 cfg 重新拉一遍元数据 if dict_data and dict_data.get("data_dictionary"):
from .step1_data_dict import run_step1 # 复用 Step 1 产出(Wave 调度器已保证 Step 1 先完成)
dict_data = run_step1(cfg, log=lambda lvl, msg, **_kw: log(lvl, msg, step="4") if log else None) columns = dict_data["data_dictionary"]
columns = dict_data["data_dictionary"] table_summary = dict_data.get("table_summary", [])
table_summary = dict_data["table_summary"] if log:
log("INFO", " · 复用 Step 1 数据字典", step="4")
else:
# 兜底:没拿到 Step 1 产出时回退自取(独立跑 Step 4 时用得到)
from .step1_data_dict import run_step1
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.get("table_summary", [])
if log:
log("INFO", " · 独立拉取数据字典(未复用 Step 1)", step="4")
by_table: dict[str, list[dict]] = defaultdict(list) by_table: dict[str, list[dict]] = defaultdict(list)
for row in columns: for row in columns:
......
...@@ -94,10 +94,14 @@ ...@@ -94,10 +94,14 @@
<el-col :span="12"> <el-col :span="12">
<el-form-item label="执行步骤"> <el-form-item label="执行步骤">
<el-checkbox-group v-model="form.steps"> <el-checkbox-group v-model="form.steps">
<el-checkbox v-for="s in stepsList" :key="s.num" :label="s.num"> <el-checkbox v-for="s in stepsList" :key="s.num" :label="s.num" :disabled="s.required">
Step {{ s.num }}: {{ s.title }} <span>
Step {{ s.num }}: {{ s.title }}
<el-tag v-if="s.required" size="small" type="danger" effect="dark" style="margin-left: 4px; font-size: 11px; height: 18px; line-height: 18px;">必选</el-tag>
</span>
</el-checkbox> </el-checkbox>
</el-checkbox-group> </el-checkbox-group>
<div class="form-hint">Step 1 为数据基础,必选且不可取消;其他步骤可并行执行</div>
</el-form-item> </el-form-item>
</el-col> </el-col>
<el-col :span="12"> <el-col :span="12">
...@@ -639,6 +643,12 @@ ...@@ -639,6 +643,12 @@
submitting.value = true; submitting.value = true;
jobResult.value = null; jobResult.value = null;
logs.value = []; logs.value = [];
// 兜底:确保必选步骤(Step 1)始终被勾选
for (const s of stepsList.value) {
if (s.required && !form.steps.includes(s.num)) {
form.steps.push(s.num);
}
}
try { try {
const r = await fetch('/api/jobs', { const r = await fetch('/api/jobs', {
method: 'POST', method: 'POST',
......
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