临床 Python 进阶路线图
⌕ /
路线图 › 第五阶段 · Agent 核心能力

第 19 章 · 任务规划与调度

本章目标:把"一个目标"变成一张可执行、可并行、可重试的任务图, 并知道在 pandas 世界里并行到底值不值得做。 配套模块:clinic/agent_planner.py 配套案例:cases/case10_TLF生成流水线.py


19.1 先看没有规划会怎样

同一个任务:"把这批数据里的问题找出来",两种 Agent 的差别:

文本
ReAct(无规划)                      Plan-and-Execute(有规划)

想 → 看 dm 结构                      【规划】
想 → 看 adsl 结构                    1. 列出可用的域
想 → 看 adae 结构                    2. 逐域做结构检查
想 → 看 vs 结构                      3. 跨域做一致性检查
想 → 跑 dm 的 QC                     4. 按严重度汇总
想 → 跑 adsl 的 QC                  【执行】
想 → 跑 adae 的 QC                   步骤 1 ✓  步骤 2 ✓
想 → 跑 vs 的 QC                     步骤 3 ✓  步骤 4 ✓
想 → 汇总                            【检查清单】每个域都做了吗?
                                     → 域清单:dm/adsl/adae/vs → 4/4 ✓

8 次模型调用                          3 次模型调用(1 规划 + 1 执行 + 1 汇总)
漏一个域?不容易发现                  漏一个域?清单对不上,立刻发现

规划买到的是"完整性可验证" —— 而临床工作里,这正是最值钱的那部分。

🔥 一句话:ReAct 解决"我不知道该做什么", 规划解决"我担心我漏做了什么"。


19.2 任务分解:拆到多细才算对

粒度是规划里最容易搞错的地方。两种错误都很常见:

拆得太细 拆得太粗
"读 dm 的前 100 行"、"读 dm 的列名"、"统计 dm 行数" 各算一个任务 "完成所有数据检查"算一个任务
每个任务都要重新装载数据 → I/O 重复 10 次 无法并行、无法重试、无法定位失败
任务图 40 个节点,模型自己都看晕 失败时只能整块重跑

判定标准(实用):

一个任务 = 一次工具调用能完成的事。 更细 —— 浪费;更粗 —— 失去了调度价值。

按这个标准,"检查 ADSL 的必填变量"是一个任务(对应一次 run_qc_checks(dataset="adsl", checks=["required_vars"])), 而"检查 ADSL"要拆,因为它包含多个检查项。

MECE:不重不漏

临床统计的程序员对这个词不陌生 —— 分析人群的定义就是要 MECE (Mutually Exclusive, Collectively Exhaustive)。任务分解遵循同一原则:

Python
# ❌ 有重叠:两个检查都会跑到 missing_rate
tasks = [
    Task("结构检查", checks=["required_vars", "missing_rate"]),
    Task("完整性检查", checks=["missing_rate", "date_pairs"]),
]

# ❌ 有遗漏:ADSL 有 5 项必查,只写了 3 项
tasks = [
    Task("检查 ADSL", checks=["required_vars", "missing_rate", "key_unique"]),
]

# ✅ MECE:每一项检查恰好出现在一个任务里
ALL_CHECKS = ["required_vars", "missing_rate", "key_unique",
              "date_pairs", "codelist", "range"]
tasks = [Task(f"检查 {c}", checks=[c]) for c in ALL_CHECKS]

⚠️ 不重不漏要用代码校验,不能靠肉眼。 一个 20 行的断言就够: python covered = [c for t in tasks for c in t.checks] assert len(covered) == len(set(covered)), "有重复检查项" assert set(covered) == set(ALL_CHECKS), f"遗漏:{set(ALL_CHECKS) - set(covered)}" 这就是把"规划质量"变成了"可测试的断言"。


19.3 Plan-and-Execute 三件套

文本
┌──────────┐    ┌──────────┐    ┌──────────┐
│  规划器   │───▶│  执行器   │───▶│ 重规划器  │
│ Planner  │    │ Executor │    │ Replan   │
└──────────┘    └──────────┘    └────┬─────┘
     ▲                               │ 有任务失败/结果异常
     └───────────────────────────────┘

任务对象本身是普通数据 —— 注意它没有 func 字段:

Python
# clinic/agent_planner.py(节选)
@dataclass
class Task:
    """一个可调度的任务。

    刻意**不持有函数引用** —— 只描述"要做什么",
    由执行器通过工具注册表去调用。这样任务图可以序列化、可以重放。
    """
    id: str
    title: str
    tool: str                         # 要调用的工具名
    args: dict = field(default_factory=dict)
    depends_on: tuple[str, ...] = ()  # 依赖的任务 id
    optional: bool = False            # 失败是否阻塞整条流水线
    status: str = "pending"           # pending/running/done/failed/skipped
    result: Any = None
    error: str | None = None
    elapsed: float = 0.0

为什么不在 Task 里放函数引用:一旦放了,任务图就没法 JSON 序列化, 断点恢复、离线重放、把计划发给同事审阅 —— 全都没了。 这是第 18 章"状态外置"在规划层的同一条纪律。


19.4 依赖 DAG 与拓扑排序

为什么必须要 DAG

TLF 生成的依赖关系天然是一张图:

文本
        adsl 衍生
       ┌────┴────┬──────────┐
       ▼         ▼          ▼
    表1 人口学  表2 AE    表3 实验室
       │         │          │
       └────┬────┴──────────┘
            ▼
        合并导出 RTF/PDF
  • 有依赖就必须后执行 → 拓扑排序
  • 没有依赖的可以并行 → 表1/表2/表3 同时跑
  • 有环就是设计错了 → 必须报错,而不是"跑到栈溢出"

实现:Kahn 算法(20 行)

Python
def topological_order(tasks: list[Task]) -> list[Task]:
    """Kahn 拓扑排序。有环时抛 PlanCycleError。"""
    by_id = {t.id: t for t in tasks}
    indeg = {t.id: len(t.depends_on) for t in tasks}
    children: dict[str, list[str]] = {t.id: [] for t in tasks}

    for t in tasks:
        for dep in t.depends_on:
            if dep not in by_id:
                raise PlanError(f"任务 {t.id} 依赖了不存在的任务 {dep}")
            children[dep].append(t.id)

    queue = [tid for tid, d in indeg.items() if d == 0]
    order: list[str] = []
    while queue:
        tid = queue.pop(0)
        order.append(tid)
        for ch in children[tid]:
            indeg[ch] -= 1
            if indeg[ch] == 0:
                queue.append(ch)

    if len(order) != len(tasks):
        stuck = sorted(set(indeg) - set(order))
        raise PlanCycleError(f"任务图存在循环依赖:{stuck}")
    return [by_id[tid] for tid in order]

🔥 "有环就报错"这件事比想象中重要。真正的环往往藏在数据里: 表 3 依赖表 2 的衍生结果,而表 2 又需要表 3 输出的一个统计量 —— 这种环在写代码时看不出来,只有在执行前跑一次拓扑排序才会暴露。 报错的成本是 1 秒,跑起来再发现的成本是半天。

分层(Level)并行

拓扑排序给的是线性顺序,但调度的价值在于知道哪些可以同时跑。 再走一步,算出每一层的任务:

Python
def levels(tasks: list[Task]) -> list[list[Task]]:
    """把任务按"最早可执行轮次"分层。同一层内的任务互不依赖,可以并行。"""
    by_id = {t.id: t for t in tasks}
    depth: dict[str, int] = {}

    def d(tid: str) -> int:
        if tid not in depth:
            depth[tid] = 0 if not by_id[tid].depends_on else \
                1 + max(d(x) for x in by_id[tid].depends_on)
        return depth[tid]

    out: dict[int, list[Task]] = {}
    for t in tasks:
        out.setdefault(d(t.id), []).append(t)
    return [out[k] for k in sorted(out)]

上面那张 TLF 图算出来就是 3 层,第 2 层的三张表可以一起跑。


19.5 并行执行:pandas 世界里的一盆冷水

先看结论,再看原因:

任务类型 单线程够用吗 并行手段 实际加速
读大批 XPT / CSV(I/O 密集) 慢 ThreadPoolExecutor 明显(等待期间让出 CPU)
pandas groupby / merge(C 层) 慢 ProcessPoolExecutor 明显,但要传路径不传数据
纯 Python 循环逐行处理 慢 ProcessPoolExecutor 明显 —— 而且先考虑改写掉这个循环
调 LLM API(网络等待) 慢 ThreadPoolExecutor 非常明显(90% 时间在等网络)
已向量化的 numpy 运算 快 通常不用 无(且更慢)

三个必须知道的坑

坑 1:GIL 不是绝对的,但也别指望多线程救 pandas

pandas 的 C 层操作会释放 GIL,所以多线程有一点点效果; 但 Python 层的调度开销会吃掉大部分收益。判断方法只有一个:自己测。

Python
# 别信任何人的结论,包括本书的 —— 在你的数据上量一次
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

def bench(fn, items, workers=4):
    t0 = time.perf_counter(); [fn(x) for x in items]           # 串行
    serial = time.perf_counter() - t0

    t0 = time.perf_counter()
    with ThreadPoolExecutor(workers) as ex: list(ex.map(fn, items))
    thread = time.perf_counter() - t0

    print(f"串行 {serial:6.2f}s  多线程 {thread:6.2f}s  "
          f"加速 {serial / thread:.2f}x")

坑 2:多进程不能传 DataFrame

ProcessPoolExecutor 会把参数序列化(pickle)后发给子进程。 传一个 24MB 的 DataFrame,光是序列化就要几秒,而且内存翻倍。

Python
# ❌ 反例:把数据传过去
with ProcessPoolExecutor(4) as ex:
    ex.map(process, [df1, df2, df3, df4])       # 每个 df 都会被 pickle 一遍

# ✅ 正例:把路径传过去,让子进程自己读
with ProcessPoolExecutor(4) as ex:
    ex.map(process_file, ["data/samples/adsl.csv",
                          "data/samples/adae.csv", ...])

这也是为什么工具层应该接收路径/数据集名而不是 DataFrame (见第 20 章的工具契约)。

坑 3:并发度不是越大越好

临床数据处理的瓶颈常常是磁盘和内存。8 个进程同时读 24MB 的 XPT, 可能比 4 个还慢(还要算上换页)。经验起点:min(4, CPU 核数),然后实测调。

并发执行器

Python
def execute_plan(plan: Plan, registry, policy,
                 max_workers: int = 4,
                 on_task: Callable[[Task], None] | None = None) -> Plan:
    """按层执行:同层并行,层间串行。"""
    plan = copy.deepcopy(plan)

    for layer in levels(plan.tasks):
        ready = [t for t in layer if all_deps_ok(t, plan)]
        if not ready:
            continue
        if len(ready) == 1 or max_workers <= 1:
            for t in ready:
                _run_one(t, registry, policy, on_task)
        else:
            with ThreadPoolExecutor(max_workers=min(max_workers, len(ready))) as ex:
                list(ex.map(lambda t: _run_one(t, registry, policy, on_task), ready))
    return plan

用线程池而不是进程池,是因为工具调用绝大多数时间在等 I/O (读文件、调 API)。需要进程级并行时再单独处理,不要一上来就上进程。


19.6 失败处理:分错误类型,不分心情

不同的失败要用不同的策略。在规划层做这个判断,不要在工具里做:

错误类型 例子 策略
validation 参数写错了(LLM 的锅) 不重试,把结构化错误回给模型让它改
not_found 数据集不存在 不重试,标记 skipped,继续后面的
transient 网络超时、临时限流 指数退避重试 2–3 次
permission 命中白名单外 不重试,直接终止并告警
internal 代码 bug 不重试,记录完整堆栈
Python
RETRYABLE = {"transient"}

def _run_one(task: Task, registry, policy, on_task=None) -> None:
    for attempt in range(1, policy.max_retries + 1):
        t0 = time.perf_counter()
        res = registry.execute(task.tool, task.args, policy)
        task.elapsed = time.perf_counter() - t0

        if res.ok:
            task.status, task.result = "done", res.data
            break

        if res.error_type in RETRYABLE and attempt < policy.max_retries:
            time.sleep(policy.backoff_base ** attempt)     # 指数退避 2s, 4s, 8s
            continue

        task.status, task.error = "failed", res.error
        break

    if on_task:
        on_task(task)

部分失败必须显式汇总

这是最容易被忽略、后果最严重的一点:

Python
def summarize(plan: Plan) -> str:
    done = [t for t in plan.tasks if t.status == "done"]
    failed = [t for t in plan.tasks if t.status == "failed"]
    skipped = [t for t in plan.tasks if t.status == "skipped"]

    if failed:
        # ⚠️ 绝不能说"检查完成" —— 必须说清楚哪几项没做
        return (f"完成 {len(done)}/{len(plan.tasks)} 项;"
                f"**{len(failed)} 项失败**({', '.join(t.title for t in failed)});"
                f"{len(skipped)} 项跳过。结论仅覆盖已完成部分。")
    return f"全部 {len(done)} 项完成。"

⚠️ 临床场景的硬要求:报告里出现"部分失败"时, 必须在结论的第一行就说清楚,不能让读者以为结果完整。 一个"看起来完整但实际漏了 3 项检查"的报告,比一个"明确写着 3 项没跑"的报告 危险得多 —— 前者会被当成完整证据使用。


19.7 什么时候该停下来问人(HITL)

Agent 自动化程度越高,越需要明确的人工确认点。判断标准是: 这个动作做错了,能不能撤销?

Python
@dataclass
class PlanPolicy:
    max_workers: int = 4
    max_retries: int = 2
    backoff_base: float = 2.0
    hitl_before: tuple[str, ...] = ("write_report", "write_xpt", "send_summary")
    """命中这些工具时暂停,等待人工确认后才继续。"""

实际效果:流水线跑到"写最终报告"这一步会停下来,把已完成的部分结果 摆在你面前,你确认后才落盘。这比"全自动跑完然后你发现结论不对"好得多。


19.8 完整案例:TLF 生成流水线

cases/case10_TLF生成流水线.py 把上面所有零件拼起来:

文本
【规划】把人月的工作排成 8 个任务
   ├── T1 派生 ADSL           (依赖:无)
   ├── T2 生成表1 人口学       (依赖:T1)
   ├── T3 生成表2 AE 汇总      (依赖:T1)
   ├── T4 生成表3 实验室移位   (依赖:T1)
   ├── T5 生成表4 生命体征     (依赖:T1)
   ├── T6 汇总质量检查         (依赖:无)
   ├── T7 合并导出             (依赖:T2,T3,T4,T5,T6)
   └── T8 写审计日志           (依赖:T7)

【拓扑分层】
   L0: T1, T6          ← 两个独立起点,并行
   L1: T2, T3, T4, T5  ← 四张表并行
   L2: T7
   L3: T8

【执行】4 并发 · 失败按类型处理 · 命中 write_report 时暂停等确认
【汇总】"完成 8/8,用时 3.4s,无失败项"

运行:

Shell
python cases/case10_TLF生成流水线.py                    # 4 并发
python cases/case10_TLF生成流水线.py --workers 1        # 对照:串行
python cases/case10_TLF生成流水线.py --fail t3          # 故意让一个任务失败,看汇总
python cases/case10_TLF生成流水线.py --dry-run          # 只打印计划,不执行

19.9 与 SAS 的对照

SAS 里做同样的事,通常是这么写的:

SAS
/* SAS 的"任务调度"通常就是按顺序跑 */
%include "01_adsl.sas";
%include "02_table1.sas";
%include "03_table2.sas";
%include "04_table3.sas";
%include "05_export.sas";

差异在哪:

维度 SAS 顺序 %include DAG 调度
依赖关系 靠文件顺序隐含表达 显式声明在 depends_on 里
并行 基本没有(或手工开多个会话) 同层自动并行
失败处理 后面继续跑(可能拿旧数据算) 依赖失败 → 下游标记 skipped,不拿脏数据
断点续跑 手工注释掉跑过的部分 状态序列化,--resume
漏跑 只能靠人肉核对日志 清单断言 + 汇总报告

🔥 最值得注意的一行是"失败处理": SAS 脚本里一个 %include 失败了,后面往往还会继续跑, 用上一次的中间数据集算出"看起来正常"的结果。 DAG 调度把这件事变成显式的:依赖没成功,下游就不跑,并且汇总里明写"跳过"。


19.10 常见误区

误区 后果 正确做法
任务里放函数引用 计划无法序列化/重放 只放 tool 名 + args
不做环检测 死循环或栈溢出 执行前跑拓扑排序
把所有任务都并行 抢内存、抢磁盘,反而更慢 按层并行,并发度实测
多进程传 DataFrame 序列化开销吃掉全部收益 传路径
失败就整个中止 前面 20 分钟白跑 分类处理,optional 任务可跳过
静默跳过失败项 报告看起来完整但实际缺失 汇总第一行就说清失败项
计划一次定死 中途数据变了还在按老计划跑 留重规划入口

19.11 本章小结

  1. 规划买到的是"完整性可验证" —— 清单对不上能立刻发现
  2. 粒度 = 一次工具调用 —— 更细浪费,更粗失去调度价值
  3. 任务只描述"做什么",不持有函数 —— 保证可序列化与可重放
  4. 拓扑排序 + 分层 —— 有环就报错,同层就并行
  5. 并行要实测 —— pandas 世界里的加速比和直觉差别很大
  6. 部分失败必须显式汇总 —— 这是临床场景的红线

下一章:工具契约怎么设计,才能让 LLM 少犯错、让错误可修复。


19.12 动手练习

  1. 跑 case10 的 --workers 1 与 --workers 4 两种模式, 记录用时。换做 --workers 8 再试一次 —— 观察是否变慢,想想为什么。
  2. 给 case10 增加一个任务 T9,让它依赖 T3 和 T4, 然后故意写成依赖 T9 自己 —— 确认 PlanCycleError 正确报出。
  3. 把 levels() 的结果画成 Mermaid 图打印出来(提示:用 graph TD 语法)。
  4. 在 _run_one 里加一个"同一任务连续失败 3 次则跳过"的规则, 并让 summarize() 在汇总里体现出来。
  5. 思考题:如果监管要求"每次生成的 TLF 必须逐字节一致", 并行会不会破坏这个要求?该怎么处理?