Commit 3e07d7ac authored by Data Governance Dev's avatar Data Governance Dev

fix(orchestrator): 依赖不在计划中时 WARN 跳过,而不是 RuntimeError

上一个提交引入了 wave 调度,如果用户勾选了有依赖的
Step 但没勾该依赖(如选 1/3/4/5 但不选 2),
调度器会在 ready=[] 时抛 RuntimeError('调度死锁'),
整个 job 标失败,日志提示看起来很严重但其实只是用户配置问题。

修复:
- 在进入 wave 调度前加预处理,把'依赖不在计划里'
  的 Step 拎出来 WARN 跳过(沿用非关键失败的处理路径)
- 把 wave 循环里 ready=[] 时的 RuntimeError 降级为 ERROR +
  break(防御性保留,理论上预处理后不应该再发生)
- 调度计划日志里多加一行 '因依赖缺失跳过=...' 字段,
  用户 / 排查时可以一眼看出被跳了哪个

验证(steps=[1,3,4,5] 此次失败案例):
  调整前:Job 失败,错误 '调度死锁: 剩余 [3] 但无依赖被满足'
  调整后:Step 3 因缺 [2] 被 WARN 跳过,Job 正常跑完 1/4/5
parent 506152ce
...@@ -176,9 +176,28 @@ async def run_governance_workflow( ...@@ -176,9 +176,28 @@ async def run_governance_workflow(
has_step8 = 8 in steps has_step8 = 8 in steps
main_steps = sorted(n for n in steps if n != 8) main_steps = sorted(n for n in steps if n != 8)
# ── 依赖完整性检查 ──
# 如果某个 Step 的前置依赖不在计划里,会造成 wave 死锁
# 此处先把这类 Step 拎出来 WARN 跳过(非关键),不要让它阻塞整个流程
planned_set = set(main_steps)
keep = []
skipped_for_deps = []
for n in main_steps:
missing = STEP_DEPENDENCIES.get(n, set()) - planned_set
if missing:
skipped_for_deps.append((n, sorted(missing)))
else:
keep.append(n)
if skipped_for_deps:
log("WARN",
f"以下步骤因缺少前置依赖被自动跳过(避免调度死锁): "
+ ", ".join(f"Step {n}(需要 {deps})" for n, deps in skipped_for_deps))
main_steps = keep
log("INFO", log("INFO",
f"调度计划: 必选={sorted(required_steps)}, " f"调度计划: 必选={sorted(required_steps)}, "
f"主步骤={main_steps}, 是否报告={has_step8}") f"主步骤={main_steps}, 是否报告={has_step8}, "
f"因依赖缺失跳过={skipped_for_deps or '无'}")
# ── Wave 调度器 ── # ── Wave 调度器 ──
# 按依赖图分波并发;每波内的 Step 通过 asyncio.gather 并行跑 # 按依赖图分波并发;每波内的 Step 通过 asyncio.gather 并行跑
...@@ -196,8 +215,11 @@ async def run_governance_workflow( ...@@ -196,8 +215,11 @@ async def run_governance_workflow(
and n not in failed_non_critical and n not in failed_non_critical
) )
if not ready: if not ready:
raise RuntimeError( # 理论上前置依赖检查后,这里不该发生;留作防御
f"调度死锁: 剩余 {sorted(remaining)} 但无依赖被满足") log("ERROR",
f"调度异常: 剩余 {sorted(remaining)} 已无可 ready 步骤"
f"(依赖图可能成环或数据异常),跳过剩余步骤")
break
wave_idx += 1 wave_idx += 1
log("INFO", f"──── Wave {wave_idx}: 并行执行 {ready} ────") log("INFO", f"──── Wave {wave_idx}: 并行执行 {ready} ────")
......
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